Anne-Marie Kermarrec

dblp:86/676 · DBLP profile ↗
← Back
179ranked-venue papers
28as first author
27since 2021 · last 2026
0000-0001-8187-724XORCID · verified

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

Systems, architecture and hardware · 62 · 11 first-author · 7 since 2021Computer networks · 30 · 4 first-author · 1 since 2021Security and privacy · 28 · 2 first-author · 5 since 2021Databases, data management, data science and information retrieval · 17 · 1 first-author · 4 since 2021Software engineering, systems software and programming languages · 15 · 2 first-author · 2 since 2021Artificial intelligence and machine learning · 7 · 1 first-author · 6 since 2021Theory of computation · 5 · 1 first-authorGraphics, computer vision, multimedia, augmented reality and games · 3 · 1 first-author · 1 since 2021Applied, interdisciplinary, general and emerging computing · 2 · 1 since 2021Human-computer interaction and ubiquitous computing · 1
YearPublicationVenuePosition
2026 Efficient Federated Search for Retrieval-Augmented Generation Using Lightweight Routing
Akash Balasaheb Dhasade, Rachid Guerraoui, Anne-Marie Kermarrec, Diana Petrescu, Rafael Pires 0001, Mathis Randl, Martijn de Vos
DAIS3
2026 HarMoEny: Efficient Inference of MoE Models
Zachary Doucet, Rishi Sharma 0001, Martijn de Vos, Rafael Pires 0001, Anne-Marie Kermarrec, Oana Balmau
IPDPS5
2025 Efficient Pyramidal Analysis of Gigapixel Images on a Decentralized Modest Computer Cluster
Marie Reinbigler, Rishi Sharma 0001, Rafael Pires 0001, Elisabeth Brunet, Anne-Marie Kermarrec, Catalin I. Fetita
Euro-Par (3)5
2025 Robust ML Auditing using Prior Knowledge
abstract
Among the many technical challenges to enforcing AI regulations, one crucial yet underexplored problem is the risk of audit manipulation. This manipulation occurs when a platform deliberately alters its answers to a regulator to pass an audit without modifying its answers to other users. In this paper, we introduce a novel approach to manipulation-proof auditing by taking into account the auditor's prior knowledge of the task solved by the platform. We first demonstrate that regulators must not rely on public priors (e.g. a public dataset), as platforms could easily fool the auditor in such cases. We then formally establish the conditions under which an auditor can prevent audit manipulations using prior knowledge about the ground truth. Finally, our experiments with two standard datasets illustrate the maximum level of unfairness a platform can hide before being detected as malicious. Our formalization and generalization of manipulation-proof auditing with a prior opens up new research directions for more robust fairness audits.
Jade Garcia Bourrée, Augustin Godinot, Sayan Biswas, Anne-Marie Kermarrec, Erwan Le Merrer, Gilles Trédan, Martijn de Vos, Milos Vujasinovic
ICML4
2025 Leveraging Approximate Caching for Faster Retrieval-Augmented Generation
Shai Bergman, Anne-Marie Kermarrec, Diana Petrescu, Rafael Pires 0001, Mathis Randl, Martijn de Vos, Ji Zhang 0035
Middleware2
2025 Robust Fingerprinting of Graphs With Fing
abstract
Graphs have become fundamental for carrying invaluable insights into numerous scientific disciplines. Controlling if they are further shared and modified is essential when sharing such graphs. This control is typically achieved using digital watermarking by embedding identification information in the graph structure. In this paper, we propose the first approach to fingerprinting graphs by associating a characteristic signature of these graphs that can be extracted later as proof of ownership. This work provides the same guarantees as watermarking while avoiding the need to modify the graph, instead by exporting the fingerprint to an external timestamped database. We present the novel fingerprinting scheme Fing. Fing relies on the Factor-r Sum Subsets problem to create a digital fingerprint. This problem is$N P$-hard, so it is easy to create and extract for the graph originator while being intractable for an attacker. We provide an analysis of the robustness of FING facing a wide range of attacks that aim at removing or extracting the fingerprint. Finally, we empirically show FING's scalability. A fingerprint can be created in around four minutes on a single core for 10 million node graphs and is robust against attacks removing thousands of edges, for instance.
Odysseas Drosis, Jade Garcia Bourrée, Anne-Marie Kermarrec, Erwan Le Merrer, Othmane Safsafi
SRDS3
2025 Boosting Asynchronous Decentralized Learning with Model Fragmentation
Sayan Biswas, Anne-Marie Kermarrec, Alexis Marouani, Rafael Pires 0001, Rishi Sharma 0001, Martijn de Vos
WWW2
2025 Noiseless Privacy-Preserving Decentralized Learning
abstract
Decentralized learning (DL) enables collaborative learning without a server and without training data leaving the users' devices. However, the models shared in DL can still be used to infer training data. Conventional defenses such as differential privacy and secure aggregation fall short in effectively safeguarding user privacy in DL, either sacrificing model utility or efficiency. We introduce Shatter, a novel DL approach in which nodes create virtual nodes (VNs) to disseminate chunks of their full model on their behalf. This enhances privacy by (i) preventing attackers from collecting full models from other nodes, and (ii) hiding the identity of the original node that produced a given model chunk. We theoretically prove the convergence of Shatter and provide a formal analysis demonstrating how Shatter reduces the efficacy of attacks compared to when exchanging full models between nodes. We evaluate the convergence and attack resilience of Shatter with existing DL algorithms, with heterogeneous datasets, and against three standard privacy attacks. Our evaluation shows that Shatter not only renders these privacy attacks infeasible when each node operates 16 VNs but also exhibits a positive impact on model utility compared to standard DL. In summary, Shatter enhances the privacy of DL while maintaining the utility and efficiency of the model.
Sayan Biswas, Mathieu Even, Anne-Marie Kermarrec, Laurent Massoulié, Rafael Pires 0001, Rishi Sharma 0001, Martijn de Vos
Proc. Priv. Enhancing Technol.3
2025 Low-Cost Privacy-Preserving Decentralized Learning
abstract
Decentralized learning (DL) is an emerging paradigm of collaborative machine learning that enables nodes in a network to train models collectively without sharing their raw data or relying on a central server. This paper introduces Zip-DL, a privacy-aware DL algorithm that leverages correlated noise to achieve robust privacy against local adversaries while ensuring efficient convergence at low communication costs. By progressively neutralizing the noise added during distributed averaging, Zip-DL combines strong privacy guarantees with high model accuracy. Its design requires only one communication round per gradient descent iteration, significantly reducing communication overhead compared to competitors. We establish theoretical bounds on both convergence speed and privacy guarantees. Moreover, extensive experiments demonstrating Zip-DL's practical applicability make it outperform state-of-the-art methods in the accuracy vs. vulnerability trade-off. Specifically, Zip-DL (i) reduces membership-inference attack success rates by up to 35% compared to baseline DL, (ii) decreases attack efficacy by up to 13% compared to competitors offering similar utility, and (iii) achieves up to 59% higher accuracy to completely nullify a basic attack scenario, compared to a state-of-the-art privacy-preserving approach under the same threat model. These results position Zip-DL as a practical and efficient solution for privacy-preserving decentralized learning in real-world applications.
Sayan Biswas, Davide Frey, Romaric Gaudel, Anne-Marie Kermarrec, Dimitri Lerévérend, Rafael Pires 0001, Rishi Sharma 0001, François Taïani
Proc. Priv. Enhancing Technol.4
2025 Boosting Resource-Constrained Federated Learning Systems With Guessed Updates
abstract
Federated learning (FL) enables a set of client devices to collaboratively train a model without sharing raw data. This process, though, operates under the constrained computation and communication resources of edge devices. These constraints combined with systems heterogeneity force some participating clients to perform fewer local updates than expected by the server, thus slowing down convergence. Exhaustive tuning of hyperparameters in FL, furthermore, can be resource-intensive, without which the convergence is adversely affected. In this work, we propose GEL, the guess and learn algorithm. GEL enables constrained edge devices to perform additional learning through guessed updates on top of gradient-based steps. These guesses aregradientless, i.e., participating clients leverage themfor free. Our generic guessing algorithm (i) can be flexibly combined with several state-of-the-art algorithms includingFedProx + GeL,FedNova,FedYogiorScaleFL; and (ii) achieves significantly improved performance when the learning rates are not best tuned. We conduct extensive experiments and show that GEL can boost empirical convergence by up to 40% in resourceconstrained networks while relieving the need for exhaustive learning rate tuning.
Mohamed Yassine Boukhari, Akash Balasaheb Dhasade, Anne-Marie Kermarrec, Rafael Pires 0001, Othmane Safsafi, Rishi Sharma 0001
IEEE Trans. Parallel Distributed Syst.3
2024 Fairness Auditing with Multi-Agent Collaboration
abstract
Existing work in fairness auditing assumes that each audit is performed independently. In this paper, we consider multiple agents working together, each auditing the same platform for different tasks. Agents have two levers: their collaboration strategy, with or without coordination beforehand, and their strategy for sampling appropriate data points. We theoretically compare the interplay of these levers. Our main findings are that (i) collaboration is generally beneficial for accurate audits, (ii) basic sampling methods often prove to be effective, and (iii) counter-intuitively, extensive coordination on queries often deteriorates audits accuracy as the number of agents increases. Experiments on three large datasets confirm our theoretical results. Our findings motivate collaboration during fairness audits of platforms that use ML models for decision-making.
Martijn de Vos, Akash Balasaheb Dhasade, Jade Garcia Bourrée, Anne-Marie Kermarrec, Erwan Le Merrer, Benoît Rottembourg, Gilles Trédan
ECAI4
2024 QuickDrop: Efficient Federated Unlearning via Synthetic Data Generation
abstract
Federated Unlearning (FU) aims to delete specific training data from an ML model trained using Federated Learning (FL). However, existing FU methods suffer from inefficiencies due to the high costs associated with gradient recomputation and storage. This paper presents QuickDrop, an original and efficient FU approach designed to overcome these limitations. During model training, each client uses QuickDrop to generate a compact synthetic dataset, serving as a compressed representation of the gradient information utilized during training. This synthetic dataset facilitates fast gradient approximation, allowing rapid downstream unlearning at minimal storage cost. To unlearn some knowledge from the trained model, QuickDrop clients execute stochastic gradient ascent with samples from the synthetic datasets instead of the training dataset. The tiny volume of synthetic data significantly reduces computational overhead compared to conventional FU methods. Evaluations with three standard datasets and five baselines show that, with comparable accuracy guarantees, QuickDrop reduces the unlearning duration by 463× compared to retraining the model from scratch and 65 -- 218× compared to FU baselines. QuickDrop supports both class- and client-level unlearning, multiple unlearning requests, and relearning of previously erased data.
Akash Balasaheb Dhasade, Yaohong Ding, Song Guo 0001, Anne-Marie Kermarrec, Martijn de Vos, Leijie Wu
Middleware4
2024 Revisiting Ensembling in One-Shot Federated Learning
abstract
Federated Learning (FL) is an appealing approach to training machine learning models without sharing raw data. However, standard FL algorithms are iterative and thus induce a significant communication cost. One-Shot FL (OFL) trades the iterative exchange of models between clients and the server with a single round of communication, thereby saving substantially on communication costs. Not surprisingly, OFL exhibits a performance gap in terms of accuracy with respect to FL, especially under high data heterogeneity. We introduce Fens, a novel federated ensembling scheme that approaches the accuracy of FL with the communication efficiency of OFL. Learning in Fens proceeds in two phases: first, clients train models locally and send them to the server, similar to OFL; second, clients collaboratively train a lightweight prediction aggregator model using FL. We showcase the effectiveness of Fens through exhaustive experiments spanning several datasets and heterogeneity levels. In the particular case of heterogeneously distributed CIFAR-10 dataset, Fens achieves up to a $26.9$% higher accuracy over SOTA OFL, being only $3.1$% lower than FL. At the same time, Fens incurs at most $4.3\times$ more communication than OFL, whereas FL is at least $10.9\times$ more communication-intensive than Fens.
Youssef Allouah, Akash Balasaheb Dhasade, Rachid Guerraoui, Nirupam Gupta, Anne-Marie Kermarrec, Rafael Pinot, Rafael Pires 0001, Rishi Sharma 0001
NeurIPS5
2024 Peerswap: A Peer-Sampler with Randomness Guarantees
abstract
The ability of a peer-to-peer (P2P) system to effectively host decentralized applications often relies on the availability of a peer-sampling service, which provides each participant with a random sample of other peers. Despite the practical effectiveness of existing peer samplers, their ability to produce random samples within a reasonable time frame remains poorly understood from a theoretical standpoint. This paper contributes to bridging this gap by introducing PeersWap,a peer-sampling protocol with provable randomness guarantees. We establish execution time bounds for PeerSwap, demonstrating its ability to scale effectively with the network size. We prove that PeerSwap maintains the fixed structure of the communication graph while allowing sequential peer position swaps within this graph. We do so by showing that PeerSwap is a specific instance of an interchange process, a renowned model for particle movement analysis. Leveraging this mapping, we derive execution time bounds, expressed as a function of the network size n. Depending on the network structure, this time can be as low as a polylogarithmic function of n, highlighting the efficiency of PeerSwap. We implement PeerSwap and conduct numerical evaluations using regular graphs with varying connectivity and containing up to 32768(215) peers. Our evaluation demonstrates that PeerSwap quickly provides peers with uniform random samples of other peers.
Rachid Guerraoui, Anne-Marie Kermarrec, Anastasiia Kucherenko, Rafael Pinot, Martijn de Vos
SRDS2
2024 Rethinking Personalized Client Collaboration in Federated Learning
abstract
Federated Learning (FL) has gained considerable attention recently, as it allows clients to cooperatively train a global machine learning model without sharing raw data. However, its performance can be compromised due to the high heterogeneity in clients' local data distributions, commonly known as Non-IID (non-independent and identically distributed). Moreover, collaboration among highly dissimilar clients exacerbates this performance degradation. Personalized FL seeks to mitigate this by enabling clients to collaborate primarily with others who have similar data characteristics, thereby producing personalized models. We noticed that existing methods for assessing model similarity often do not capture the genuine relevance of client domains. In response, our paper enhances personalized client collaboration in FL by introducing a metric for domain relevance between clients. Specifically, to facilitate optimal coalition formation, we measure the marginal contributions of client models using coalition game theory, providing a more accurate representation of potential client domain relevance within the FL privacy-preserving framework. Based on this metric, we then adjust each client's coalition membership and implement a personalized FL aggregation algorithm that is robust to Non-IID data domain. We provide a theoretical analysis of the algorithm's convergence and generalization capabilities. Our extensive evaluations on multiple datasets, including MNIST, Fashion-MNIST, CIFAR-10, and CIFAR-100, and under varying Non-IID data distributions (Pathological and Dirichlet), demonstrate that our personalized collaboration approach consistently outperforms contemporary benchmarks in terms of accuracy for individual clients.
Leijie Wu, Song Guo 0001, Yaohong Ding, Wenchao Xu 0001, Yufeng Zhan, Anne-Marie Kermarrec
IEEE Trans. Mob. Comput.7
2023 Refined Convergence and Topology Learning for Decentralized SGD with Heterogeneous Data
abstract
One of the key challenges in decentralized and federated learning is to design algorithms that efficiently deal with highly heterogeneous data distributions across agents. In this paper, we revisit the analysis of Decentralized Stochastic Gradient Descent algorithm (D-SGD) under data heterogeneity. We exhibit the key role played by a new quantity, called neighborhood heterogeneity, on the convergence rate of D-SGD. By coupling the communication topology and the heterogeneity, our analysis sheds light on the poorly understood interplay between these two concepts. We then argue that neighborhood heterogeneity provides a natural criterion to learn data-dependent topologies that reduce (and can even eliminate) the otherwise detrimental effect of data heterogeneity on the convergence time of D-SGD. For the important case of classification with label skew, we formulate the problem of learning such a good topology as a tractable optimization problem that we solve with a Frank-Wolfe algorithm. As illustrated over a set of simulated and real-world experiments, our approach provides a principled way to design a sparse topology that balances the convergence speed and the per-iteration communication costs of D-SGD under data heterogeneity.
Batiste Le Bars, Aurélien Bellet, Marc Tommasi, Erick Lavoie, Anne-Marie Kermarrec
AISTATS5
2023 Get More for Less in Decentralized Learning Systems
abstract
Decentralized learning (DL) systems have been gaining popularity because they avoid raw data sharing by communicating only model parameters, hence preserving data confidentiality. However, the large size of deep neural networks poses a significant challenge for decentralized training, since each node needs to exchange gigabytes of data, overloading the network. In this paper, we address this challenge with Jwins, a communication-efficient and fully decentralized learning system that shares only a subset of parameters through sparsification. Jwins uses wavelet transform to limit the information loss due to sparsification and a randomized communication cut-off that reduces communication usage without damaging the performance of trained models. We demonstrate empirically with 96 DL nodes on non-IID datasets that Jwins can achieve similar accuracies to full-sharing DL while sending up to 64% fewer bytes. Additionally, on low communication budgets, Jwins outperforms the state-of-the-art communication-efficient DL algorithm Choco-SGD by up to 4x in terms of network savings and time.
Akash Balasaheb Dhasade, Anne-Marie Kermarrec, Rafael Pires 0001, Rishi Sharma 0001, Milos Vujasinovic, Jeffrey Wigger
ICDCS2
2023 Epidemic Learning: Boosting Decentralized Learning with Randomized Communication
abstract
We present Epidemic Learning (EL), a simple yet powerful decentralized learning (DL) algorithm that leverages changing communication topologies to achieve faster model convergence compared to conventional DL approaches. At each round of EL, each node sends its model updates to a random sample of $s$ other nodes (in a system of $n$ nodes). We provide an extensive theoretical analysis of EL, demonstrating that its changing topology culminates in superior convergence properties compared to the state-of-the-art (static and dynamic) topologies. Considering smooth non-convex loss functions, the number of transient iterations for EL, i.e., the rounds required to achieve asymptotic linear speedup, is in $O(n^3/s^2)$ which outperforms the best-known bound $O(n^3)$ by a factor of $s^2$, indicating the benefit of randomized communication for DL. We empirically evaluate EL in a 96-node network and compare its performance with state-of-the-art DL approaches. Our results illustrate that EL converges up to $ 1.7\times$ quicker than baseline DL algorithms and attains $2.2 $\% higher accuracy for the same communication volume.
Martijn de Vos, Sadegh Farhadkhani, Rachid Guerraoui, Anne-Marie Kermarrec, Rafael Pires 0001, Rishi Sharma 0001
NeurIPS4
2023 On the Inherent Anonymity of Gossiping
abstract
Detecting the source of a gossip is a critical issue, related to identifying patient zero in an epidemic, or the origin of a rumor in a social network. Although it is widely acknowledged that random and local gossip communications make source identification difficult, there exists no general quantification of the level of anonymity provided to the source. This paper presents a principled method based on $\varepsilon$-differential privacy to analyze the inherent source anonymity of gossiping for a large class of graphs. First, we quantify the fundamental limit of source anonymity any gossip protocol can guarantee in an arbitrary communication graph. In particular, our result indicates that when the graph has poor connectivity, no gossip protocol can guarantee any meaningful level of differential privacy. This prompted us to further analyze graphs with controlled connectivity. We prove on these graphs that a large class of gossip protocols, namely cobra walks, offers tangible differential privacy guarantees to the source. In doing so, we introduce an original proof technique based on the reduction of a gossip protocol to what we call a random walk with probabilistic die out. This proof technique is of independent interest to the gossip community and readily extends to other protocols inherited from the security community, such as the Dandelion protocol. Interestingly, our tight analysis precisely captures the trade-off between dissemination time of a gossip protocol and its source anonymity.
Rachid Guerraoui, Anne-Marie Kermarrec, Anastasiia Kucherenko, Rafael Pinot, Sasha Voitovych
DISC2
2023 GoldFinger: Fast & Approximate Jaccard for Efficient KNN Graph Constructions
abstract
We proposeGoldFinger, a newcompactandfast-to-computebinary representation of datasets to approximate Jaccard's index. We illustrate the effectiveness of GoldFinger on the emblematic big data problem of K-Nearest-Neighbor (KNN) graph construction and show that GoldFinger can drastically accelerate a large range of existing KNN algorithms with little to no overhead. As a side effect, we also show that the compact representation of the data protects users’ privacyfor freeby providingk-anonymity andl-diversity. Our extensive evaluation of the resulting approach on several realistic datasets shows that our approach reduces computation times by up to 78.9% compared to raw data while only incurring a negligible to moderate loss in terms of KNN quality. We also show that GoldFinger can be applied to KNN queries (a widely-used search technique) and delivers speedups of up to$\times 3.55$over one of the most efficient approaches to this problem.
Rachid Guerraoui, Anne-Marie Kermarrec, Guilhem Niot, Olivier Ruas, François Taïani
IEEE Trans. Knowl. Data Eng.2
2022 TEE-based decentralized recommender systems: The raw data sharing redemption
abstract
Recommenders are central in many applications today. The most effective recommendation schemes, such as those based on collaborative filtering (CF), exploit similarities between user profiles to make recommendations, but potentially expose private data. Federated learning and decentralized learning systems address this by letting the data stay on user's machines to preserve privacy: each user performs the training on local data and only the model parameters are shared. However, sharing the model parameters across the network may still yield privacy breaches. In this paper, we present Rex, the first enclave-based decentralized CF recommender. Rex exploits Trusted execution environments (TEE), such as Intel software guard extensions (SGX), that provide shielded environments within the processor to improve convergence while preserving privacy. Firstly, Rex enables raw data sharing, which ultimately speeds up convergence and reduces the network load. Secondly, Rex fully preserves privacy. We analyze the impact of raw data sharing in both deep neural network (DNN) and matrix factorization (MF) recommenders and showcase the benefits of trusted environments in a full-fledged implementation of Rex. Our experimental results demonstrate that through raw data sharing, Rex significantly decreases the training time by 18.3 x and the network load by 2 orders of magnitude over standard decentralized approaches that share only parameters, while fully protecting privacy by leveraging trustworthy hardware enclaves with very little overhead.
Akash Balasaheb Dhasade, Nevena Dresevic, Anne-Marie Kermarrec, Rafael Pires 0001
IPDPS3
2022 The Universal Gossip Fighter
abstract
The notion of adversary is a staple of distributed computing. An adversary typically models “hostile” assumptions about the underlying distributed environment, e.g., a network that can drop messages, an operating system that can delay processes or an attacker that can hack machines. So far, the goal of distributed computing researchers has mainly been to develop a distributed algorithm that can face a given adversary, the abstraction characterizing worst-case scenarios. This paper initiates the study of the somehow opposite approach. Given a distributed algorithm, the adversary is the abstraction we seek to implement. More specifically, we consider the problem of controlling the spread of messages in a large-scale system, conveying the practical motivation of limiting the dissemination of fake news or viruses. Essentially, we assume a general class of gossip protocols, called all-to-all gossip protocols, and devise a practical method to hinder the dissemination. We present the Universal Gossip Fighter (UGF). Just like classical adversaries in distributed computing, UGF can observe the status of a dissemination and decide to stop some processes or delay some messages. The originality of UGF lies in the fact that it is universal, i.e., it applies to any all-to-all gossip protocol. We show that any gossip protocol attacked by UGF ends up exhibiting a quadratic message complexity (in the total number of processes) if it achieves sublinear time of dissemination. We also show that if a gossip protocol aims to achieve a message complexity$\alpha$times smaller than quadratic, then the time complexity rises exponentially in relation to$\alpha$. We convey the practical relevance of our theoretical findings by implementing UGF and conducting a set of empirical experiments that confirm some of our results.
Anastasiia Gorbunova, Rachid Guerraoui, Anne-Marie Kermarrec, Anastasiia Kucherenko, Rafael Pinot
IPDPS3
2022 Frugal Decentralized Learning
abstract
Machine learning is currently shifting from a centralized paradigm to decentralized ones where machine learning models are trained collaboratively. In fully decentralized learning algorithms, data remains where it was produced, models are trained locally and only model parameters are exchanged among participating entities along an arbitrary network topology and aggregated over time until convergence. Not only this limits the cost of exchanging data but also exploits the growing capabilities of users' devices while mitigating privacy and confidentiality concerns. Such systems are significantly challenged by a potential high-level of heterogeneity both at the system level as participants may have differing capabilities of (e.g., computing power, memory and network connectivity) as well as data heterogeneity (a.k.a non-IIDness). The adoption of fully decentralized learning systems requires designing frugal systems that limit communication, energy and yet ensure convergences. Several avenues are promising from adapting the network topologies to compensate for data heterogeneity to exploiting the high levels of redundancy, both in data and computations, of ML algorithms to limit data and model sharing in such systems.
Anne-Marie Kermarrec
IPDPS1
2022 D-Cliques: Compensating for Data Heterogeneity with Topology in Decentralized Federated Learning
abstract
The convergence speed of machine learning models trained with Federated Learning is significantly affected by heterogeneous data partitions, even more so in a fully decentralized setting without a central server. In this paper, we show that the impact of label distribution skew, an important type of data heterogeneity, can be significantly reduced by carefully designing the underlying communication topology. We present D-Cliques, a novel topology that reduces gradient bias by grouping nodes in sparsely interconnected cliques such that the label distribution in a clique is representative of the global label distribution. We also show how to adapt the updates of decentralized SGD to obtain unbiased gradients and implement an effective momentum with D-Cliques. Our extensive empirical evaluation on MNIST and CIFAR10 validates our design and demonstrates that our approach achieves similar convergence speed as a fully-connected topology, while providing a significant reduction in the number of edges and messages. In a 1000-node topology, D-Cliques require 98% less edges and 96% less total messages, with further possible gains using a small-world topology across cliques.
Aurélien Bellet, Anne-Marie Kermarrec, Erick Lavoie
SRDS2
2022 FLeet: Online Federated Learning via Staleness Awareness and Performance Prediction
abstract
Federated learning (FL) is very appealing for its privacy benefits: essentially, a global model is trained with updates computed on mobile devices while keeping the data of users local. Standard FL infrastructures are however designed to have no energy or performance impact on mobile devices, and are therefore not suitable for applications that require frequent ( online ) model updates, such as news recommenders. This article presents FLeet , the first Online FL system, acting as a middleware between the Android operating system and the machine learning application. FLeet combines the privacy of Standard FL with the precision of online learning thanks to two core components: (1) I-Prof , a new lightweight profiler that predicts and controls the impact of learning tasks on mobile devices, and (2) AdaSGD , a new adaptive learning algorithm that is resilient to delayed updates. Our extensive evaluation shows that Online FL, as implemented by FLeet , can deliver a 2.3× quality boost compared to Standard FL while only consuming 0.036% of the battery per day. I-Prof can accurately control the impact of learning tasks by improving the prediction accuracy by up to 3.6× in terms of computation time, and by up to 19× in terms of energy. AdaSGD outperforms alternative FL approaches by 18.4% in terms of convergence speed on heterogeneous data.
Georgios Damaskinos, Rachid Guerraoui, Anne-Marie Kermarrec, Vlad Nitu, Rhicheek Patra, François Taïani
ACM Trans. Intell. Syst. Technol.3
2021 Cluster-and-Conquer: When Randomness Meets Graph Locality
abstract
K-Nearest-Neighbors (KNN) graphs are central to many emblematic data mining and machine-learning applications. Some of the most efficient KNN graph algorithms are incremental and local: they start from a random graph, which they incrementally improve by traversing neighbors-of-neighbors links. Unfortunately, the initial random graph exhibits a poor graph locality, leading to many unnecessary similarity computations. In this paper, we remove this drawback with Cluster-and-Conquer (C2for short). Cluster-and-Conquer boosts the starting configuration of greedy algorithms thanks to a novel lightweight clustering mechanism, dubbed FastRandomHash. FastRandomHash leverages randomness and recursion to pre-cluster similar nodes at a very low cost. Our extensive evaluation on real datasets shows that Cluster-and-Conquer significantly outperforms existing approaches, including LSH, yielding speed-ups of up to ×4.42 and even improving the KNN quality.
George Giakkoupis, Anne-Marie Kermarrec, Olivier Ruas, François Taïani
ICDE2
2021 Quicker ADC : Unlocking the Hidden Potential of Product Quantization With SIMD
abstract
Efficient Nearest Neighbor (NN) search in high-dimensional spaces is a foundation of many multimedia retrieval systems. A common approach is to rely on Product Quantization, which allows the storage of large vector databases in memory and efficient distance computations. Yet, implementations of nearest neighbor search with Product Quantization have their performance limited by the many memory accesses they perform. Following this observation, André et al. proposed Quick ADC with up to 6× faster implementations of PQ m×4 product quantizers (PQ) leveraging specific SIMD instructions. Quicker ADC is a generalization of Quick ADC not limited to PQ m×4 codes and supporting AVX-512, the latest revision of SIMD instruction set. In doing so, Quicker ADC faces the challenge of using efficiently 5,6 and 7-bit shuffles that do not align to computer bytes or words. To this end, we introduce (i) irregular product quantizers combining sub-quantizers of different granularity and (ii) split tables allowing lookup tables larger than registers. We evaluate Quicker ADC with multiple indexes including Inverted Multi-Indexes and IVF HNSW and show that it outperforms the reference optimized implementations (i.e., FAISS and polysemous codes) for numerous configurations. Finally, we release an open-source fork of FAISS enhanced with Quicker ADC.
Fabien André, Anne-Marie Kermarrec, Nicolas Le Scouarnec
IEEE Trans. Pattern Anal. Mach. Intell.2
2020 FLeet: Online Federated Learning via Staleness Awareness and Performance Prediction
abstract
Federated Learning (FL) is very appealing for its privacy benefits: essentially, a global model is trained with updates computed on mobile devices while keeping the data of users local. Standard FL infrastructures are however designed to have no energy or performance impact on mobile devices, and are therefore not suitable for applications that require frequent (online) model updates, such as news recommenders.
Georgios Damaskinos, Rachid Guerraoui, Anne-Marie Kermarrec, Vlad Nitu, Rhicheek Patra, François Taïani
Middleware3
2020 FeGAN: Scaling Distributed GANs
abstract
Existing approaches to distribute Generative Adversarial Networks (GANs) either (i) fail to scale for they typically put the two components of a GAN (the generator and the discriminator) on different machines, inducing significant communication overhead, or (ii) they face GAN training specific issues, exacerbated by distribution.
Rachid Guerraoui, Arsany Guirguis, Anne-Marie Kermarrec, Erwan Le Merrer
Middleware3
2020 Smaller, Faster & Lighter KNN Graph Constructions
abstract
We propose GoldFinger, a new compact and fast-to-compute binary representation of datasets to approximate Jaccard’s index. We illustrate the effectiveness of GoldFinger on the emblematic big data problem of K-Nearest-Neighbor (KNN) graph construction and show that GoldFinger can drastically accelerate a large range of existing KNN algorithms with little to no overhead. As a side effect, we also show that the compact representation of the data protects users’ privacy for free by providing k-anonymity and l-diversity. Our extensive evaluation of the resulting approach on several realistic datasets shows that our approach delivers speedups of up to 78.9% compared to the use of raw data while only incurring a negligible to moderate loss in terms of KNN quality. To convey the practical value of such a scheme, we apply it to item recommendation and show that the loss in recommendation quality is negligible.
Rachid Guerraoui, Anne-Marie Kermarrec, Olivier Ruas, François Taïani
WWW2
2019 Fingerprinting Big Data: The Case of KNN Graph Construction
abstract
We propose fingerprinting, a new technique that consists in constructing compact, fast-to-compute and privacy-preserving binary representations of datasets. We illustrate the effectiveness of our approach on the emblematic big data problem of K-Nearest-Neighbor (KNN) graph construction and show that fingerprinting can drastically accelerate a large range of existing KNN algorithms, while efficiently obfuscating the original data, with little to no overhead. Our extensive evaluation of the resulting approach (dubbed GoldFinger) on several realistic datasets shows that our approach delivers speedups of up to 78.9% compared to the use of raw data while only incurring a negligible to moderate loss in terms of KNN quality.
Rachid Guerraoui, Anne-Marie Kermarrec, Olivier Ruas, François Taïani
ICDE2
2018 Collaborative Filtering Under a Sybil Attack: Similarity Metrics do Matter!
abstract
Recommendation systems help users identify interesting content, but they also open new privacy threats. In this paper, we deeply analyze the effect of a Sybil attack that tries to infer information on users from a user-based collaborative-filtering recommendation systems. We discuss the impact of different similarity metrics used to identity users with similar tastes in the trade-off between recommendation quality and privacy. Finally, we propose and evaluate a novel similarity metric that combines the best of both worlds: a high recommendation quality with a low prediction accuracy for the attacker. Our results, on a state-of-the-art recommendation framework and on real datasets show that existing similarity metrics exhibit a wide range of behaviors in the presence of Sybil attacks, while our new similarity metric consistently achieves the best trade-off while outperforming state-of-the-art solutions.
Antoine Boutet, Florestan De Moor, Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Antoine Rault
DSN5
2018 Nobody Cares if You Liked Star Wars: KNN Graph Construction on the Cheap
Anne-Marie Kermarrec, Olivier Ruas, François Taïani
Euro-Par1
2017 I Know Nothing about You But Here is What You Might Like
abstract
Recommenders widely use collaborative filtering schemes. These schemes, however, threaten privacy as user profiles are made available to the service provider hosting the recommender and can even be guessed by curious users who analyze the recommendations. Users can encrypt their profiles to hide them from the service provider and add noise to make them difficult to guess. These precautionary measures hamper latency and recommendation quality. In this paper, we present a novel recommender, X-REC, enabling an effective collaborative filtering scheme to ensure the privacy of users against the service provider (system-level privacy) or other users (user-level privacy). X-REC builds on two underlying services: X-HE, an encryption scheme designed for recommenders, and X-NN, a neighborhood selection protocol over encrypted profiles. We leverage uniform sampling to ensure differential privacy against curious users. Our extensive evaluation demonstrates that X-REC provides (1) recommendation quality similar to non-private recommenders, and (2) significant latency improvement over privacy-aware alternatives.
Rachid Guerraoui, Anne-Marie Kermarrec, Rhicheek Patra, Mahammad Valiyev, Jingjing Wang 0007
DSN2
2017 Agar: A Caching System for Erasure-Coded Data
abstract
Erasure coding is an established data protection mechanism. It provides high resiliency with low storage overhead, which makes it very attractive to storage systems developers. Unfortunately, when used in a distributed setting, erasure coding hampers a storage system's performance, because it requires clients to contact several, possibly remote sites to retrieve their data. This has hindered the adoption of erasure coding in practice, limiting its use to cold, archival data. Recent research showed that it is feasible to use erasure coding for hot data as well, thus opening new perspectives for improving erasure-coded storage systems. In this paper, we address the problem of minimizing access latency in erasure-coded storage. We propose Agar-a novel caching system tailored for erasure-coded content. Agar optimizes the contents of the cache based on live information regarding data popularity and access latency to different data storage sites. Our system adapts a dynamic programming algorithm to optimize the choice of data blocks that are cached, using an approach akin to "Knapsack" algorithms. We compare Agar to the classical Least Recently Used and Least Frequently Used cache eviction policies, while varying the amount of data cached between a data chunk and a whole replica of the object. We show that Agar can achieve 16% to 41% lower latency than systems that use classical caching policies.
Raluca Halalai, Pascal Felber, Anne-Marie Kermarrec, François Taïani
ICDCS3
2017 Accelerated Nearest Neighbor Search with Quick ADC
abstract
Efficient Nearest Neighbor (NN) search in high-dimensional spaces is a foundation of many multimedia retrieval systems. Because it offers low responses times, Product Quantization (PQ) is a popular solution. PQ compresses high-dimensional vectors into short codes using several sub-quantizers, which enables in-RAM storage of large databases. This allows fast answers to NN queries, without accessing the SSD or HDD. The key feature of PQ is that it can compute distances between short codes and high-dimensional vectors using cache-resident lookup tables. The efficiency of this technique, named Asymmetric Distance Computation (ADC), remains limited because it performs many cache accesses.
Fabien André, Anne-Marie Kermarrec, Nicolas Le Scouarnec
ICMR2
2017 The Utility and Privacy Effects of a Click
abstract
Recommenders are becoming one of the main ways to navigate the Internet. They recommend appropriate items to users based on their clicks, i.e., likes, ratings, purchases, etc. These clicks are key to providing relevant recommendations and, in this sense, have a significant utility. Since clicks reflect the preferences of users, they also raise privacy concerns. At first glance, there seems to be an inherent trade-off between the utility and privacy effects of a click. Nevertheless, a closer look reveals that the situation is more subtle: some clicks do improve utility without compromising privacy, whereas others decrease utility while hampering privacy.
Rachid Guerraoui, Anne-Marie Kermarrec, Mahsa Taziki
SIGIR2
2017 Recommenders: from the Lab to the Wild (Keynote Talk)
abstract
Recommenders are ubiquitous on the Internet today: they tell you which book to read, which movie you should watch, predict your next holiday destination, give you advices on restaurants and hotels, they are even responsible for the posts that you see on your favorite social media and potentially greatly influence your friendship on social networks. While many approaches exist, collaborative filtering is one of the most popular approaches to build online recommenders that provide users with content that matches their interest. Interestingly, the very notion of users can be general and span actual humans or software applications. Recommenders come with many challenges beyond the quality of the recommendations. One of the most prominent ones is their ability to scale to a large number of users and a growing volume of data to provide real-time recommendations introducing many system challenges. Another challenge is related to privacy awareness: while recommenders rely on the very fact that users give away information about themselves, this potentially raises some privacy concerns. In this talk, I will focus on the challenges associated to building efficient, scalable and privacy-aware recommenders.
Anne-Marie Kermarrec
DISC1
2017 Heterogeneous Recommendations: What You Might Like To Read After Watching Interstellar
abstract
Recommenders, as widely implemented nowadays by major e-commerce players like Netflix or Amazon, use collaborative filtering to suggest the most relevant items to their users. Clearly, the effectiveness of recommenders depends on the data they can exploit, i.e., the feedback of users conveying their preferences, typically based on their past ratings. As of today, most recommenders are homogeneous in the sense that they utilize one specific application at a time. In short, Alice will only get recommended a movie if she has been rating movies. But what if she has been only rating books and would like to get recommendations for a movie? Clearly, the multiplicity of web applications is calling for heterogeneous recommenders that could utilize ratings in one application to provide recommendations in another one. This paper presents X-M ap , a heterogeneous recommender. X-M ap leverages meta-paths between heterogeneous items over several application domains, based on users who rated across these domains. These meta-paths are then used in X-M ap to generate, for every user, a profile ( AlterEgo ) in a domain where the user might not have rated any item yet. Not surprisingly, leveraging meta-paths poses non-trivial issues of (a) meta-path-based inter-item similarity , in order to enable accurate predictions, (b) scalability , given the amount of computation required, as well as (c) privacy , given the need to aggregate information across multiple applications. We show in this paper how X-M ap addresses the above-mentioned issues to achieve accuracy, scalability and differential privacy. In short, X-M ap weights the meta-paths based on several factors to compute inter-item similarities, and ensures scalability through a layer-based pruning technique. X-M ap guarantees differential privacy using an exponential scheme that leverages the meta-path-based similarities while determining the probability of item selection to construct the AlterEgos. We present an exhaustive experimental evaluation of X-M ap using real traces from Amazon. We show that, in terms of accuracy, X-M ap outperforms alternative heterogeneous recommenders and, in terms of throughput, X-M ap achieves a linear speedup with an increasing number of machines.
Rachid Guerraoui, Anne-Marie Kermarrec, Tao Lin 0004, Rhicheek Patra
Proc. VLDB Endow.2
2016 ProteusTM: Abstraction Meets Performance in Transactional Memory
abstract
The Transactional Memory (TM) paradigm promises to greatly simplify the development of concurrent applications. This led, over the years, to the creation of a plethora of TM implementations delivering wide ranges of performance across workloads. Yet, no universal implementation fits each and every workload. In fact, the best TM in a given workload can reveal to be disastrous for another one. This forces developers to face the complex task of tuning TM implementations, which significantly hampers their wide adoption. In this paper, we address the challenge of automatically identifying the best TM implementation for a given workload. Our proposed system, ProteusTM, hides behind the TM interface a large library of implementations. Underneath, it leverages a novel multi-dimensional online optimization scheme, combining two popular learning techniques: Collaborative Filtering and Bayesian Optimization.
Diego Didona, Nuno Diegues, Anne-Marie Kermarrec, Rachid Guerraoui, Ricardo Neves, Paolo Romano 0002
ASPLOS3
2016 Bounds on the Voter Model in Dynamic Networks
abstract
In the voter model, each node of a graph has an opinion, and in every round each node chooses independently a random neighbour and adopts its opinion. We are interested in the consensus time, which is the first point in time where all nodes have the same opinion. We consider dynamic graphs in which the edges are rewired in every round (by an adversary) giving rise to the graph sequence G_1, G_2, ..., where we assume that G_i has conductance at least phi_i. We assume that the degrees of nodes don't change over time as one can show that the consensus time can become super-exponential otherwise. In the case of a sequence of d-regular graphs, we obtain asymptotically tight results. Even for some static graphs, such as the cycle, our results improve the state of the art. Here we show that the expected number of rounds until all nodes have the same opinion is bounded by O(m/(d_{min}*phi)), for any graph with m edges, conductance phi, and degrees at least d_{min}. In addition, we consider a biased dynamic voter model, where each opinion i is associated with a probability P_i, and when a node chooses a neighbour with that opinion, it adopts opinion i with probability P_i (otherwise the node keeps its current opinion). We show for any regular dynamic graph, that if there is an epsilon > 0 difference between the highest and second highest opinion probabilities, and at least Omega(log(n)) nodes have initially the opinion with the highest probability, then all nodes adopt w.h.p. that opinion. We obtain a bound on the convergence time, which becomes O(log(n)/phi) for static graphs.
Petra Berenbrink, George Giakkoupis, Anne-Marie Kermarrec, Frederik Mallmann-Trenn
ICALP3
2016 Being prepared in a sparse world: The case of KNN graph construction
abstract
K-Nearest-Neighbor (KNN) graphs have emerged as a fundamental building block of many on-line services providing recommendation, similarity search and classification. Constructing a KNN graph rapidly and accurately is, however, a computationally intensive task. As data volumes keep growing, speed and the ability to scale out are becoming critical factors when deploying a KNN algorithm. In this work, we present KIFF, a generic, fast and scalable KNN graph construction algorithm. KIFF directly exploits the bipartite nature of most datasets to which KNN algorithms are applied. This simple but powerful strategy drastically limits the computational cost required to rapidly converge to an accurate KNN solution, especially for sparse datasets. Our evaluation on a representative range of datasets show that KIFF provides, on average, a speed-up factor of 14 against recent state-of-the art solutions while improving the quality of the KNN approximation by 18%.
Antoine Boutet, Anne-Marie Kermarrec, Nupur Mittal, François Taïani
ICDE2
2016 Atum: Scalable Group Communication Using Volatile Groups
Rachid Guerraoui, Anne-Marie Kermarrec, Matej Pavlovic, Dragos-Adrian Seredinschi
Middleware2
2016 Distributed Slicing in Dynamic Systems
abstract
Peer to peer (P2P) systems have moved from application specific architectures to a generic service oriented design philosophy. This raised interesting problems in connection with providing useful P2P middleware services capable of dealing with resource assignment and management in a large-scale, heterogeneous and unreliable environment. The slicing problem consists of partitioning a P2P network into$k$groups (slices) of a given portion of the network nodes that share similar resource values. As the network is large and dynamic this partitioning is continuously updated without any node knowing the network size. In this paper, we propose the first algorithm to solve the slicing problem. We introduce the metric of slice disorder and show that the existing ordering algorithm cannot nullify this disorder. We propose a new algorithm that speeds up the existing ordering algorithm but that suffers from the same inaccuracy. Then, we propose another algorithm based on ranking that is provably convergent under reasonable assumptions. In particular, we notice experimentally that ordering algorithms suffer from resource-correlated churn while the ranking algorithm can cope with it. These algorithms are proved viable theoretically and experimentally.
Antonio Fernández 0001, Vincent Gramoli, Ernesto Jiménez, Anne-Marie Kermarrec, Michel Raynal
IEEE Trans. Parallel Distributed Syst.4
2015 Similitude: Decentralised Adaptation in Large-Scale P2P Recommenders
Davide Frey, Anne-Marie Kermarrec, Christopher Maddock, Andreas Mauthe, Pierre-Louis Roman, François Taïani
DAIS2
2015 Cheap and Cheerful: Trading Speed and Quality for Scalable Social-Recommenders
Anne-Marie Kermarrec, François Taïani, Juan Manuel Tirado
DAIS1
2015 Hide & Share: Landmark-Based Similarity for Private KNN Computation
abstract
Computing k-nearest-neighbor graphs constitutes a fundamental operation in a variety of data-mining applications. As a prominent example, user-based collaborative-filtering provides recommendations by identifying the items appreciated by the closest neighbors of a target user. As this kind of applications evolve, they will require KNN algorithms to operate on more and more sensitive data. This has prompted researchers to propose decentralized peer-to-peer KNN solutions that avoid concentrating all information in the hands of one central organization. Unfortunately, such decentralized solutions remain vulnerable to malicious peers that attempt to collect and exploit information on participating users. In this paper, we seek to overcome this limitation by proposing H&S (Hide & Share), a novel landmark-based similarity mechanism for decentralized KNN computation. Landmarks allow users (and the associated peers) to estimate how close they lay to one another without disclosing their individual profiles. We evaluate H&S in the context of a user-based collaborative-filtering recommender with publicly available traces from existing recommendation systems. We show that although landmark-based similarity does disturb similarity values (to ensure privacy), the quality of the recommendations is not as significantly hampered. We also show that the mere fact of disturbing similarity values turns out to be an asset because it prevents a malicious user from performing a profile reconstruction attack against other users, thus reinforcing users' privacy. Finally, we provide a formal privacy guarantee by computing an upper bound on the amount of information revealed by H&S about a user's profile.
Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Antoine Rault, François Taïani, Jingjing Wang 0007
DSN3
2015 CPSys: A System for Mobile Video Prefetching
abstract
Online media services are reshaping the way video content is watched. People with similar interests tend to request same content. This provides enormous potential to predict which content users are interested in. Besides, mobile devices are commonly used to watch videos which popularity is largely driven by its social success. In this paper, we design CPSys a Central Predictor System to prefetch relevant videos for each user. To fine tune our prefetching system, we rely on a large dataset collected from a large mobile carrier in Europe. The rationale of our prefetching strategy is first to form a graph and build implicit or explicit ties between similar users. On top of this graph, we propose the Most Popular and Most Recent (MPMR) policy to predict relevant videos for each user. We show that CPSys can achieve high performance with respect to the correct prediction ratio and by significantly reducing the traffic overhead. We further show that CPSys outperforms other prefetching schemes that have been presented and studied in the state of the art. At the end, we provide a proof-of-concept implementation of our prefetching system.
Ali Gouta, David Hausheer, Anne-Marie Kermarrec, Christian Koch 0003, Yannick Le Louédec, Julius Rückert
MASCOTS3
2015 Scaling Out Link Prediction with SNAPLE: 1 Billion Edges and Beyond
abstract
A growing number of organizations are seeking to analyze extra large graphs in a timely and resource-efficient manner. With some graphs containing well over a billion elements, these organizations are turning to distributed graph-computing platforms that can scale out easily in existing data-centers and clouds. Unfortunately such platforms usually impose programming models that can be ill suited to typical graph computations, fundamentally undermining their potential benefits.
Anne-Marie Kermarrec, François Taïani, Juan Manuel Tirado
Middleware1
2015 Hawk: Hybrid Datacenter Scheduling
Pamela Delgado, Florin Dinu, Anne-Marie Kermarrec, Willy Zwaenepoel
USENIX ATC3
2015 Privacy-Conscious Information Diffusion in Social Networks
abstract
We present Riposte , a distributed algorithm for disseminating information (ideas, news, opinions, or trends) in a social network. Riposte ensures that information spreads widely if and only if a large fraction of users find it interesting, and this is done in a “privacy-conscious” manner, namely without revealing the opinion of any individual user. Whenever an information item is received by a user, Riposte decides to either forward the item to all the user’s neighbors, or not to forward it to anyone. The decision is randomized and is based on the user’s (private) opinion on the item, as well as on an upper bound s on the number of user’s neighbors that have not received the item yet. In short, if the user likes the item, Riposte forwards it with probability slightly larger than 1 / s , and if not, the item is forwarded with probability slightly smaller than 1 / s . Using a comparison to branching processes, we show for a general family of random directed graphs with arbitrary out-degree sequences, that if the information item appeals to a sufficiently large (constant) fraction of users, then the item spreads to a constant fraction of the network; while if fewer users like it, the dissemination process dies out quickly. In addition, we provide extensive experimental evaluation of Riposte on topologies taken from online social networks, including Twitter and Facebook. 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.
George Giakkoupis, Rachid Guerraoui, Arnaud Jégou, Anne-Marie Kermarrec, Nupur Mittal
DISC4
2015 WebGC Gossiping on Browsers Without a Server [Live Demo/Poster]
Raziel Carvajal-Gomez, Davide Frey, Matthieu Simonin, Anne-Marie Kermarrec
WISE (2)4
2015 Cache locality is not enough: High-Performance Nearest Neighbor Search with Product Quantization Fast Scan
abstract
Nearest Neighbor (NN) search in high dimension is an important feature in many applications (e.g., image retrieval, multimedia databases). Product Quantization (PQ) is a widely used solution which offers high performance, i.e., low response time while preserving a high accuracy. PQ represents high-dimensional vectors (e.g., image descriptors) by compact codes. Hence, very large databases can be stored in memory, allowing NN queries without resorting to slow I/O operations. PQ computes distances to neighbors using cache-resident lookup tables, thus its performance remains limited by (i) the many cache accesses that the algorithm requires, and (ii) its inability to leverage SIMD instructions available on modern CPUs. In this paper, we advocate that cache locality is not sufficient for efficiency. To address these limitations, we design a novel algorithm, PQ Fast Scan, that transforms the cache-resident lookup tables into small tables, sized to fit SIMD registers. This transformation allows (i) in-register lookups in place of cache accesses and (ii) an efficient SIMD implementation. PQ Fast Scan has the exact same accuracy as PQ, while having 4 to 6 times lower response time (e.g., for 25 million vectors, scan time is reduced from 74ms to 13ms).
Fabien André, Anne-Marie Kermarrec, Nicolas Le Scouarnec
Proc. VLDB Endow.2
2015 D2P: Distance-Based Differential Privacy in Recommenders
abstract
The upsurge in the number of web users over the last two decades has resulted in a significant growth of online information. This information growth calls for recommenders that personalize the information proposed to each individual user. Nevertheless, personalization also opens major privacy concerns. This paper presents D 2 P , a novel protocol that ensures a strong form of differential privacy, which we call distance-based differential privacy, and which is particularly well suited to recommenders. D 2 P avoids revealing exact user profiles by creating altered profiles where each item is replaced with another one at some distance. We evaluate D 2 P analytically and experimentally on MovieLens and Jester datasets and compare it with other private and non-private recommenders.
Rachid Guerraoui, Anne-Marie Kermarrec, Rhicheek Patra, Mahsa Taziki
Proc. VLDB Endow.2
2014 Behave: Behavioral Cache for Web Content
Davide Frey, Mathieu Goessens, Anne-Marie Kermarrec
DAIS3
2014 Archiving cold data in warehouses with clustered network coding
abstract
Modern storage systems now typically combine plain replication and erasure codes to reliably store large amount of data in datacenters. Plain replication allows a fast access to popular data, while erasure codes, e.g., Reed-Solomon codes, provide a storage-efficient alternative for archiving less popular data. Although erasure codes are now increasingly employed in real systems, they experience high overhead during maintenance, i.e., upon failures, typically requiring files to be decoded before being encoded again to repair the encoded blocks stored at the faulty node.
Fabien André, Anne-Marie Kermarrec, Erwan Le Merrer, Nicolas Le Scouarnec, Gilles Straub, Alexandre van Kempen
EuroSys2
2014 Polystyrene: the Decentralized Data Shape That Never Dies
abstract
Decentralized topology construction protocols organize nodes along a predefined topology (e.g. a torus, ring, or hypercube). Such topologies have been used in many contexts ranging from routing and storage systems, to publish-subscribe and event dissemination. Since most topologies assume no correlation between the physical location of nodes and their positions in the topology, they do not handle catastrophic failures well, in which a whole region of the topology disappears. When this occurs, the overall shape of the system typically gets lost. This is highly problematic in applications in which overlay nodes are used to map a virtual data space, be it for routing, indexing or storage. In this paper, we propose a novel decentralized approach that maintains the initial shape of the topology even if a large (consecutive) portion of the topology fails. Our approach relies on the dynamic decoupling between physical nodes and virtual ones enabling a fast reshaping. For instance, our results show that a 51,200-node torus converges back to a full torus in only 10 rounds after 50% of the nodes have crashed. Our protocol is both simple and flexible and provides a novel form of collective survivability that goes beyond the current state of the art.
Simon Bouget, Hoel Kervadec, Anne-Marie Kermarrec, François Taïani
ICDCS3
2014 HyRec: leveraging browsers for scalable recommenders
abstract
The ever-growing amount of data available on the Internet calls for personalization. Yet, the most effective personalization schemes, such as those based on collaborative filtering (CF), are notoriously resource greedy. This paper presents HyRec, an online cost-effective scalable system for user-based CF personalization. HyRec offloads recommendation tasks onto the web browsers of users, while a server orchestrates the process and manages the relationships between user profiles.
Antoine Boutet, Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Rhicheek Patra
Middleware4
2014 Improving caching efficiency and quality of experience with CF-Dash
abstract
HTTP Adaptive Streaming (HAS) is gradually being adopted by Over The Top (OTT) content providers. In HAS, a wide range of video bitrates of the same video content are made available over the internet so that clients' players pick the video bitrate that best fit their bandwidth. Yet, this affects the performance of some major components of the video delivery chain, namely CDNs or transparent caches since several versions of the same content compete to be cached. In this context we investigate the benefits of a Cache Friendly HAS system (CF-DASH), which aims to improve the caching efficiency in mobile networks and to sustain the quality of experience of mobile clients. Firstly, we motivate our work by presenting a set of observations we made on large number of clients requesting HAS contents. Secondly we introduce the CF-Dash system and our testbed implementation. Finally, we evaluate CF-dash based on trace-driven simulations and testbed experiments. Our validation results are promising. Simulations on real HAS traffic show that we achieve a significant gain in hit-ratio that ranges from 15% up to 50%..
Zied Aouini, Mamadou Tourad Diallo, Ali Gouta, Anne-Marie Kermarrec, Yannick Le Louédec
NOSSDAV4
2014 Comparing the Predictive Capability of Social and Interest Affinity for Recommendations
Alexandra Olteanu, Anne-Marie Kermarrec, Karl Aberer
WISE (1)2
2014 Tracking freeriders in gossip-based content dissemination systems
abstract
Gossip-based protocols have proven very efficient for disseminating high-bandwidth content such as video streams in a peer-to-peer fashion. However, for the protocols to work, nodes are required to collaborate by devoting a fraction of their upload bandwidth, a scarce resource for some of them, to forward the content they receive to other nodes. Consequently, such protocols suffer from freeriding, a common phenomenon on the Internet, which consists in selfishly benefiting from the system without contributing its fair share. Due to the dynamic nature and the inherent randomness of gossip protocols and to the high scalability requirements of video streaming systems, detecting freeriders is a difficult challenge. This paper presents LiFTinG, the first protocol for detecting freeriders, including colluding ones, in gossip-based content dissemination systems with asymmetric data exchanges. In addition, LiFTinG is still able to detect freeriders when network coding, a widely used technique to improve the efficiency of content dissemination, is used. LiFTinG relies on nodes to track abnormal behavior by cross-checking the history of their previous interactions and exploits the fact that nodes pick neighbors at random to prevent colluding nodes from mutually covering up their bad actions. We present a methodology for setting the parameters of LiFTinG to their optimal value, based on a theoretical analysis and we quantify theoretically the performance of LiFTinG. We derive, based on simulations, the optimal strategy of freeriders by taking into account, through a utility function, the benefit of freeriding and the probability of being detected. In addition to these simulations, we report on the deployment of LiFTinG on PlanetLab. In a 300-node system, where a stream of 674 kbps is broadcasted, LiFTinG incurs a maximum overhead of only 8% and provides good detection results: For instance, with 10% of freeriders decreasing their contribution by up to 30%, LiFTinG detects 86% of the freeriders after only 30 s and wrongfully expels only a few honest nodes (most of them actually being buggy).
Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod, Swagatika Prusty, Aline Roumy
Comput. Networks3
2014 Performance evaluation of a peer-to-peer backup system using buffering at the edge
Anne-Marie Kermarrec, Erwan Le Merrer, Nicolas Le Scouarnec, Romaric Ludinard, Patrick Maillé, Gilles Straub, Alexandre van Kempen
Comput. Commun.1
2014 Computing in social networks
Andrei Giurgiu, Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec
Inf. Comput.4
2014 Personalizing Top-k Processing Online in a Peer-to-Peer Social Tagging Network
abstract
The rapidly increasing amount of user-generated content in social tagging systems provides a huge source of information. Yet, performing effective search in these systems is very challenging, especially when we seek the most appropriate items that match a potentially ambiguous query. Collaborative filtering-based personalization is appealing in this context, as it limits the search within a small network of participants with similar preferences. Offline personalization, which consists in maintaining, for every user, a network of similar participants based on their tagging behaviors, is effective for queries that are close to the querying user’s tagging profile but performs poorly when the queries, reflecting emerging interests, have little correlation with the querying user’s profile. We present P 2 TK 2 , the first protocol to personalize query processing in social tagging systems online. P 2 TK 2 is completely decentralized, and this design choice stems from the observation that the evolving social tagging systems naturally resemble P2P systems where users are both producers and consumers. This design exploits the power of the crowd and prevents any central authority from controlling personal information. P 2 TK 2 is gossip-based and probabilistic. It dynamically associates each user with social acquaintances sharing similar tagging behaviors. Appropriate users for answering a query are discovered at query time with the help of social acquaintances. This is achieved according to the hybrid interest of the querying user, taking into account both her tagging behavior and her query. Results are iteratively refined and returned to the querying user. We evaluate P 2 TK 2 on CiteULike and Delicious traces involving up to 50,000 users. We highlight the advantages of online personalization compared to offline personalization, as well as its efficiency, scalability, and inherent ability to cope with user departure and interest evolution in P2P systems.
Xiao Bai 0002, Rachid Guerraoui, Anne-Marie Kermarrec
ACM Trans. Internet Techn.3
2014 Convex Partitioning of Large-Scale Sensor Networks in Complex Fields: Algorithms and Applications
abstract
When a sensor network grows large, or when its topology becomes complex (e.g., containing many holes), network algorithms designed with a smaller or simpler setting in mind may be rendered rather inefficient. We propose to address this problem using a divide and conquer approach: the network is divided into convex pieces by a distributed convex partitioning protocol, using connectivity information only. A convex network partition exhibits some desirable properties that allow traditional algorithms to work to their full advantage. Based on this, we can achieve relatively high performance for an algorithm by combining algorithmic actions within individual partitions. We consider two important applications: virtual-coordinate-based geographic routing and connectivity-based localization. The former benefits from convex partition's friendliness to network embedding, which is crucial to generating accurate virtual coordinates for the nodes, while the latter leverages the fact that shortest paths are largely straight for node pairs within a convex partition. Experimental results show that the convex partition approach can significantly improve the performance of both applications in comparison with state-of-the-art solutions.
Guang Tan, Hongbo Jiang 0001, Anne-Marie Kermarrec
ACM Trans. Sens. Networks4
2013 FlexGD: A flexible force-directed model for graph drawing
abstract
We propose FlexGD, a force-directed algorithm for straightline undirected graph drawing. The algorithm strives to draw graph layouts encompassing from uniform vertex distribution to extreme structure abstraction. It is flexible for it is parameterized so that the emphasis can be put on either of the two drawing criteria. The parameter determines how much the edges are shorter than the average distance between vertices. Extending the clustering property of the LinLog model, FlexGD is efficient for cluster visualization in an adjustable level. The energy function of FlexGD is minimized through a multilevel approach, particularly designed to work in contexts where edge length distribution is not uniform. Applying FlexGD on several real datasets, we illustrate both the good quality of the layout on various topologies, and the ability of the algorithm to meet the addressed drawing criteria.
Anne-Marie Kermarrec, Afshin Moin
PacificVis1
2013 WHATSUP: A Decentralized Instant News Recommender
abstract
We present WHATSUP, a collaborative filtering system for disseminating news items in a large-scale dynamic setting with no central authority. WHATSUP constructs an implicit social network based on user profiles that express the opinions of users about the news items they receive (like-dislike). Users with similar tastes are clustered using a similarity metric reflecting long-standing and emerging (dis)interests. News items are disseminated through a novel heterogeneous gossip protocol that (1) biases the orientation of its targets towards those with similar interests, and (2) amplifies dissemination based on the level of interest in every news item. We report on an extensive evaluation of WHATSUP through (a) simulations, (b) a ModelNet emulation on a cluster, and (c) a PlanetLab deployment based on real datasets. We show that WHATSUP outperforms various alternatives in terms of accurate and complete delivery of relevant news items while preserving the fundamental advantages of standard gossip: namely, simplicity of deployment and robustness.
Antoine Boutet, Davide Frey, Rachid Guerraoui, Arnaud Jégou, Anne-Marie Kermarrec
IPDPS5
2013 HTTP Adaptive Streaming in Mobile Networks: Characteristics and Caching Opportunities
abstract
Cellular networks have witnessed the emergence of the HTTP Adaptive Streaming (HAS) as a new video delivery method. In HAS, several qualities of the same videos are made available in the network so that clients can choose the best quality that fits their bandwidth capacity. This has particular implications on caching strategies with respect to the viewing patterns and the switching behavior between video qualities. In this paper we present analysis of a real HAS dataset collected in France and provided by the country's largest mobile phone operator. Firstly, we analyse the viewing patterns of HAS contents and the distribution of the encoding bit rates requested by mobile clients. Secondly, we give an in-depth analysis of the switching pattern between video bit rates during a video session and assess the implication on the caching efficiency. We also model this switching based on empirical observations. Finally, we propose WA-LRU a new caching algorithm tailored for HAS contents and compare it to the standard LRU. Our evaluations demonstrate that WA-LRU performs better and achieves its goals.
Ali Gouta, Dohy Hong, Anne-Marie Kermarrec, Yannick Le Louédec
MASCOTS3
2013 Highly dynamic distributed computing with byzantine failures
abstract
This paper shows for the first time that distributed computing can be both reliable and efficient in an environment that is both highly dynamic and hostile. More specifically, we show how to maintain clusters of size O(log N), each containing more than two thirds of honest nodes with high probability, within a system whose size can vary polynomially with respect to its initial size. Furthermore, the communication cost induced by each node arrival or departure is polylogarithmic with respect to N, the maximal size of the system. Our clustering can be achieved despite the presence of a Byzantine adversary controlling a fraction τ ≤ 1⁄3-ε of the nodes, for some fixed constant ε > 0, independent of N. So far, such a clustering could only be performed for systems whose size can vary constantly and it was not clear whether that was at all possible for polynomial variances.
Rachid Guerraoui, Florian Huc, Anne-Marie Kermarrec
PODC3
2013 Gossip Protocols for Renaming and Sorting
George Giakkoupis, Anne-Marie Kermarrec, Philipp Woelfel
DISC2
2013 Large scale analysis of HTTP Adaptive Streaming in mobile networks
abstract
HTTP Adaptive bitrate video Streaming (HAS) is now widely adopted by Content Delivery Network Providers (CDNPs) and Telecom Operators (Telcos) to improve user Quality of Experience (QoE). In HAS, several versions of videos are made available in the network so that the quality of the video can be chosen to better fit the bandwidth capacity of users. These delivery requirements raise new challenges with respect to content caching strategies, since several versions of the content may compete to be cached. In this paper we present analysis of a real HAS dataset collected in France and provided by a mobile telecom operator involving more than 485,000 users requesting adaptive video contents through more than 8 million video sessions over a 6 week measurement period. Firstly, we propose a fine-grained definition of content popularity by exploiting the segmented nature of video streams. We also provide analysis about the behavior of clients when requesting such HAS streams. We propose novel caching policies tailored for chunk-based streaming. Then we study the relationship between the requested video bitrates and radio constraints. Finally, we study the users' patterns when selecting different bitrates of the same video content. Our findings provide useful insights that can be leveraged by the main actors of video content distribution to improve their content caching strategy for adaptive streaming contents as well as to model users' behavior in this context.
Ali Gouta, Charles Hong, Dohy Hong, Anne-Marie Kermarrec, Yannick Le Louédec
WOWMOM4
2013 Byzantine agreement with homonyms
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Anne-Marie Kermarrec, Eric Ruppert, Hung Tran-The
Distributed Comput.4
2013 Trust-aware peer sampling: Performance and privacy tradeoffs
abstract
The ability to identify people that share one’s own interests is one of the most interesting promises of the Web 2.0 driving user-centric applications such as recommendation systems or collaborative marketplaces. To be truly useful, however, information about other users also needs to be associated with some notion of trust. Consider a user wishing to sell a concert ticket. Not only must she find someone who is interested in the concert, but she must also make sure she can trust this person to pay for it. This paper addresses the need for trust in user-centric applications by proposing two novel distributed protocols that combine interest-based connections between users with explicit links obtained from social networks à-la Facebook. Both protocols build trusted multi-hop paths between users in an explicit social network supporting the creation of semantic overlays backed up by social trust. The first protocol, TAPS2, extends our previous work on TAPS (Trust-Aware Peer Sampling), by improving the ability to locate trusted nodes. Yet, it remains vulnerable to attackers wishing to learn about trust values between arbitrary pairs of users. The second protocol, PTAPS (Private TAPS), improves TAPS2 with provable privacy guarantees by preventing users from revealing their friendship links to users that are more than two hops away in the social network. In addition to proving this privacy property, we evaluate the performance of our protocols through event-based simulations, showing significant improvements over the state of the art.
Davide Frey, Arnaud Jégou, Anne-Marie Kermarrec, Michel Raynal, Julien Stainer
Theor. Comput. Sci.3
2013 Connectivity-based and anchor-free localization in large-scale 2D/3D sensor networks
abstract
A connectivity-based and anchor-free three-dimensional localization (CATL) scheme is presented for large-scale sensor networks with concave regions. It distinguishes itself from previous work with a combination of three features: (1) it works for networks in both 2D and 3D spaces, possibly containing holes or concave regions; (2) it is anchor-free and uses only connectivity information to faithfully recover the original network topology, up to scaling and rotation; (3) it does not depend on the knowledge of network boundaries, which suits it well to situations where boundaries are difficult to identify. The key idea of CATL is to discover the notch nodes , where shortest paths bend and hop-count-based distance starts to significantly deviate from the true Euclidean distance. An iterative protocol is developed that uses a notch-avoiding multilateration mechanism to localize the network. Simulations show that CATL achieves accurate localization results with a moderate per-node message cost.
Guang Tan, Hongbo Jiang 0001, Shengkai Zhang, Zhimeng Yin 0001, Anne-Marie Kermarrec
ACM Trans. Sens. Networks5
2012 Probabilistic deduplication for cluster-based storage systems
abstract
The need to backup huge quantities of data has led to the development of a number of distributed deduplication techniques that aim to reproduce the operation of centralized, single-node backup systems in a cluster-based environment. At one extreme, stateful solutions rely on indexing mechanisms to maximize deduplication. However the cost of these strategies in terms of computation and memory resources makes them unsuitable for large-scale storage systems. At the other extreme, stateless strategies store data blocks based only on their content, without taking into account previous placement decisions, thus reducing the cost but also the effectiveness of deduplication.
Davide Frey, Anne-Marie Kermarrec, Konstantinos Kloudas
SoCC2
2012 Geology: Modular Georecommendation in Gossip-Based Social Networks
abstract
Geolocated social networks, combining traditional social networking features with geolocation information, have grown tremendously over the last few years. Yet, very few works have looked at implementing geolocated social networks in a fully distributed manner, a promising avenue to handle the growing scalability challenges of these systems. In this paper, we focus on georecommendation, and show that existing decentralized recommendation mechanisms perform in fact poorly on geodata. We propose a set of novel gossip-based mechanisms to address this problem, in a modular similarity framework called GEOLOGY. The resulting platform is lightweight, efficient, and scalable, and we demonstrate its advantages in terms of recommendation quality and communication overhead on a real dataset of 15,694 users from Foursquare, a leading geolocated social network.
Jesús Carretero 0001, Florin Isaila, Anne-Marie Kermarrec, François Taïani, Juan Manuel Tirado
ICDCS3
2012 Content and geographical locality in user-generated content sharing systems
abstract
User Generated Content (UGC), such as YouTube videos, accounts for a substantial fraction of the Internet traffic. To optimize their performance, UGC services usually rely on both proactive and reactive approaches that exploit spatial and temporal locality in access patterns. Alternative types of locality are also relevant and hardly ever considered together. In this paper, we show on a large (more than 650,000 videos) YouTube dataset that content locality (induced by the related videos feature) and geographic locality, are in fact correlated. More specifically, we show how the geographic view distribution of a video can be inferred to a large extent from that of its related videos. We leverage these findings to propose a UGC storage system that proactively places videos close to the expected requests. Compared to a caching-based solution, our system decreases by 16% the number of requests served from a different country than that of the requesting user, and even in this case, the distance between the user and the server is 29% shorter on average.
Kévin Huguenin, Anne-Marie Kermarrec, Konstantinos Kloudas, François Taïani
NOSSDAV2
2012 Scalable and Secure Polling in Dynamic Distributed Networks
abstract
We consider the problem of securely conducting a poll in synchronous dynamic networks equipped with a Public Key Infrastructure (PKI). Whereas previous distributed solutions had a communication cost of O(n2) in an n nodes system, we present SPP (Secure and Private Polling), the first distributed polling protocol requiring only a communication complexity of O(n log3n), which we prove is near-optimal. Our protocol ensures perfect security against a computationally-bounded adversary, tolerates (1/2 - ϵ)n Byzantine nodes for any constant 1/2 >; ϵ >; 0 (not depending on n), and outputs the exact value of the poll with high probability. SPP is composed of two sub-protocols, which we believe to be interesting on their own: SPP-Overlay maintains a structured overlay when nodes leave or join the network, and SPP-Computation conducts the actual poll. We validate the practicality of our approach through experimental evaluations and describe briefly two possible applications of SPP: (1) an optimal Byzantine Agreement protocol whose communication complexity is Θ(n log n) and (2) a protocol solving an open question of King and Saia in the context of aggregation functions, namely on the feasibility of performing multiparty secure aggregations with a communication complexity of o(n2).
Sébastien Gambs, Rachid Guerraoui, Hamza Harkous, Florian Huc, Anne-Marie Kermarrec
SRDS5
2012 Regenerating Codes: A System Perspective
abstract
The explosion of the amount of data stored in cloud systems calls for more efficient paradigms for redundancy. While replication is widely used to ensure data availability, erasure correcting codes provide a much better trade-off between storage and availability. Regenerating codes are good candidates for they also offer low repair costs in term of network bandwidth. While they have been proven optimal, they are difficult to understand and parameterize. In this paper we provide an analysis of regenerating codes for practitioners to grasp the various trade-offs. More specifically we make two contributions: (i) we study the impact of the parameters by conducting an analysis at the level of the system, rather than at the level of a single device, (ii) we compare the computational costs of various implementations of codes and highlight the most efficient ones. Our goal is to provide system designers with concrete information to help them choose the best parameters and design for regenerating codes.
Steve Jiekak, Anne-Marie Kermarrec, Nicolas Le Scouarnec, Gilles Straub, Alexandre van Kempen
SRDS2
2012 Availability-Based Methods for Distributed Storage Systems
abstract
Distributed storage systems rely heavily on redundancy to ensure data availability as well as durability. In networked systems subject to intermittent node unavailability, the level of redundancy introduced in the system should be minimized and maintained upon failures. Repairs are well-known to be extremely bandwidth-consuming and it has been shown that, without care, they may significantly congest the system. In this paper, we propose an approach to redundancy management accounting for nodes heterogeneity with respect to availability. We show that by using the availability history of nodes, the performance of two important faces of distributed storage (replica placement and repair) can be significantly improved. Replica placement is achieved based on complementary nodes with respect to nodes availability, improving the overall data availability. Repairs can be scheduled thanks to an adaptive per-node timeout according to node availability, so as to decrease the number of repairs while reaching comparable availability. We propose practical heuristics for those two issues. We evaluate our approach through extensive simulations based on real and well-known availability traces. Results clearly show the benefits of our approach with regards to the critical trade-off between data availability, load-balancing and bandwidth consumption.
Anne-Marie Kermarrec, Erwan Le Merrer, Gilles Straub, Alexandre van Kempen
SRDS1
2012 BLIP: Non-interactive Differentially-Private Similarity Computation on Bloom filters
Mohammad Alaggan, Sébastien Gambs, Anne-Marie Kermarrec
SSS3
2012 Decentralized polling with respectable participants
Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod, Ymir Vigfusson
J. Parallel Distributed Comput.3
2012 Pulp: An adaptive gossip-based dissemination protocol for multi-source message streams
Pascal Felber, Anne-Marie Kermarrec, Lorenzo Leonini, Etienne Rivière, Spyros Voulgaris
Peer-to-Peer Netw. Appl.2
2012 Greedy Geographic Routing in Large-Scale Sensor Networks: A Minimum Network Decomposition Approach
abstract
In geographic (or geometric) routing, messages are by default routed in a greedy manner: The current node always forwards a message to its neighbor node that is closest to the destination. Despite its simplicity and general efficiency, this strategy alone does not guarantee delivery due to the existence of local minima (or dead ends). Overcoming local minima requires nodes to maintain extra nonlocal state or to use auxiliary mechanisms. We study how to facilitate greedy forwarding by using a minimum amount of such nonlocal states in topologically complex networks. Specifically, we investigate the problem of decomposing a given network into a minimum number of greedily routable components (GRCs), where greedy routing is guaranteed to work. We approach it by considering an approximate version of the problem in a continuous domain, with a central concept called the greedily routable region (GRR). A full characterization of GRR is given concerning its geometric properties and routing capability. We then develop simple approximate algorithms for the problem. These results lead to a practical routing protocol that has a routing stretch below 7 in a continuous domain, and close to 1 in several realistic network settings.
Guang Tan, Anne-Marie Kermarrec
IEEE/ACM Trans. Netw.2
2011 Distributed social graph embedding
abstract
Distributed recommender systems are becoming increasingly important for they address both scalability and the Big Brother syndrome. Link prediction is one of the core mechanism in recommender systems and relies on extracting some notion of proximity between entities in a graph. Applied to social networks, defining a proximity metric between users enable to predict potential relevant future relationships. In this paper, we propose SoCS (Social Coordinate Systems}, a fully distributed algorithm that embeds any social graph in an Euclidean space, which can easily be used to implement link prediction. To the best of our knowledge, SoCS is the first system explicitly relying on graph embedding. Inspired by recent works on non-isomorphic embeddings, the SoCS embedding preserves the community structure of the original graph, while being easy to decentralize. Nodes thus get assigned coordinates that reflect their social position. We show through experiments on real and synthetic data sets that these coordinates can be exploited for efficient link prediction.
Anne-Marie Kermarrec, Vincent Leroy 0001, Gilles Trédan
CIKM1
2011 Can Everybody Sit Closer to Their Friends Than Their Enemies?
Anne-Marie Kermarrec, Christopher Thraves
MFCS1
2011 Private Similarity Computation in Distributed Systems: From Cryptography to Differential Privacy
Mohammad Alaggan, Sébastien Gambs, Anne-Marie Kermarrec
OPODIS3
2011 Efficient peer-to-peer backup services through buffering at the edge
abstract
The availability of end devices of peer-to-peer storage and backup systems has been shown critical for usability and for system reliability in practice. This has led to the adoption of hybrid architectures composed of both peers and servers. Such architectures mask the instability of peers thus approaching the performances of client-server systems while providing scalability at a low cost. In this paper, we advocate the replacement of such servers by a cloud of residential gateways, as they are already present in users' homes, thus pushing the required stable components at the edge of the network. In our gateway-assisted system, gateways act as buffers between peers, compensating for their intrinsic instability. This enables to offload backup tasks quickly from the user's machine to the gateway, while significantly lowering the retrieval time of backed up data. We evaluate our proposal using real world traces including existing traces from Skype and Jabber as well as a trace of residential gateways for availability, and a residential broadband trace for bandwidth. Results show that the time required to backup data in the network is comparable to a server-assisted approach, while substantially improving the time to restore data, which drops from a few days to a few hours. As gateways are becoming increasingly powerful in order to enable new services, we expect such a proposal to be leveraged on a short term basis.
Serge Defrance, Anne-Marie Kermarrec, Erwan Le Merrer, Nicolas Le Scouarnec, Gilles Straub, Alexandre van Kempen
Peer-to-Peer Computing2
2011 Converging Quickly to Independent Uniform Random Topologies
abstract
The peer sampling service is a core building block for gossip protocols in peer-to-peer networks. Ideally, a peer sampling service continuously provides each peer with a sample of peers picked uniformly at random in the network. While empirical studies have shown that uniformity was achieved, analysis proposed so far assume strong restrictions on the topology of the overlay network it continuously generates. In this work, we analyze a Generic Random Peer Sampling Service (GRPS) that satisfies the desirable properties for any peer sampling service-small views, uniform sample, load balancing, and independence- and relieve strong degree connections in the nodes assumed in previous works. The main result we prove is: starting from any simple (without loops and parallel edges) directed graph with out-degree equal to c for all nodes, and recursively applying GRPS, eventually results in a random simple directed graph with out-degree equal to c for all nodes. We test empirically convergence time and independence time for GRPS. Finally, We use this empirical evaluation to show that GRPS performs better than previously presented peer sampling services.
Anne-Marie Kermarrec, Vincent Leroy 0001, Christopher Thraves
PDP1
2011 Byzantine agreement with homonyms
abstract
International audience
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Anne-Marie Kermarrec, Eric Ruppert, Hung Tran-The
PODC4
2011 Brief Announcement: A Stable and Robust Membership Protocol
Ajoy K. Datta, Anne-Marie Kermarrec, Lawrence L. Larmore, Erwan Le Merrer
SSS2
2011 Social Market: Combining Explicit and Implicit Social Networks
Davide Frey, Arnaud Jégou, Anne-Marie Kermarrec
SSS3
2011 Second order centrality: Distributed assessment of nodes criticity in complex networks
Anne-Marie Kermarrec, Erwan Le Merrer, Bruno Sericola, Gilles Trédan
Comput. Commun.1
2011 Collaborative personalized top-k processing
abstract
This article presents P4Q, a fully decentralized gossip-based protocol to personalize query processing in social tagging systems. P4Q dynamically associates each user with social acquaintances sharing similar tagging behaviors. Queries are gossiped among such acquaintances, computed on-the-fly in a collaborative, yet partitioned manner, and results are iteratively refined and returned to the querier. Analytical and experimental evaluations convey the scalability of P4Q for top- k query processing, as well its inherent ability to cope with users updating profiles and departing.
Xiao Bai 0002, Rachid Guerraoui, Anne-Marie Kermarrec, Vincent Leroy 0001
ACM Trans. Database Syst.3
2011 Analysis of Deterministic Tracking of Multiple Objects Using a Binary Sensor Network
abstract
Let consider a set of anonymous moving objects to be tracked in a binary sensor network. This article studies the problem of associating deterministically a track revealed by the sensor network with the trajectory of an unique anonymous object, namely the multiple object tracking and identification (MOTI) problem. In our model, the network is represented by a sparse connected graph where each vertex represents a binary sensor and there is an edge between two sensors if an object can pass from one sensed region to another one without activating any other sensor. The difficulty of MOTI lies in the fact that the trajectories of two or more objects can be so close that the corresponding tracks on the sensor network can no longer be distinguished (track merging), thus confusing the deterministic association between an object trajectory and a track. The article presents several results. We first show that MOTI cannot be solved on a general graph of ideal binary sensors even by an omniscient external observer if all the objects can freely move on the graph. Then we describe restrictions that can be imposed a priori either on the graph, on the object movements, or on both, to make the MOTI problem always solvable. In the absence of an omniscient observer, we show how our results can lead to the definition of distributed algorithms that are able to detect when the system is in a state where MOTI becomes unsolvable.
Yann Busnel, Leonardo Querzoni, Roberto Baldoni, Marin Bertier, Anne-Marie Kermarrec
ACM Trans. Sens. Networks5
2010 Gossiping personalized queries
abstract
International audience
Xiao Bai 0002, Marin Bertier, Rachid Guerraoui, Anne-Marie Kermarrec, Vincent Leroy 0001
EDBT4
2010 Message-Efficient Byzantine Fault-Tolerant Broadcast in a Multi-hop Wireless Sensor Network
abstract
We consider message-efficient broadcast tolerating Byzantine faults in a multi-hop wireless sensor network. Assuming a grid network where all nodes have a communication range of r, and a single neighborhood contains at most t dishonest and collision-capable (bad) nodes, each with a message budget mf, we investigate the minimum message budget m that each honest (good) node must have in order to achieve reliable broadcast. We consider three cases: (1) mfis known in advance and m is homogeneous among all good nodes; (2) mfis known in advance and m is heterogeneous among good nodes; (3) mfis unknown. For the first two cases, we present possibility results and broadcast protocols that have message costs within twice the lower bound. For the third case, we present a coding scheme that helps verify the integrity of messages at a receiving node without using any cryptographic techniques. This code leads to a reactive local broadcast primitive that has probabilistic reliability guarantees. Combined with a previously proposed scheme, it results in a broadcast protocol for t <; 1/2r(2r + 1) that guarantees reliability with high probability.
Marin Bertier, Anne-Marie Kermarrec, Guang Tan
ICDCS2
2010 LT Network Codes
abstract
Network coding has been successfully applied in large-scale content dissemination systems. While network codes provide optimal throughput, its current forms suffer from a high decoding complexity. This is an issue when applied to systems composed of nodes with low processing capabilities, such as sensor networks. In this paper, we propose a novel network coding approach based on LT codes, initially introduced in the context of erasure coding. Our coding scheme, called LTNC, fully benefits from the low complexity of belief propagation decoding. Yet, such decoding schemes are extremely sensitive to statistical properties of the code. Maintaining such properties in a fully decentralized way with only a subset of encoded data is challenging. This is precisely what the recoding algorithms of LTNC achieve. We evaluate LTNC against random linear network codes in an epidemic content-dissemination application. Results show that LTNC increases communication overhead (20\%) and convergence time (30\%) but greatly reduces the decoding complexity (99%) when compared to random linear network codes. In addition, LTNC consistently outperforms dissemination protocols without codes, thus preserving the benefit of coding.
Mary-Luc Champel, Kévin Huguenin, Anne-Marie Kermarrec, Nicolas Le Scouarnec
ICDCS3
2010 The Gossple Anonymous Social Network
Marin Bertier, Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Vincent Leroy 0001
Middleware4
2010 LiFTinG: Lightweight Freerider-Tracking in Gossip
Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod, Swagatika Prusty
Middleware3
2010 Greedy geographic routing in large-scale sensor networks: a minimum network decomposition approach
abstract
In geographic (or geometric) routing, messages are expected to route in a greedy manner: the current node always forwards a message to its neighbor node that is closest to the destination. Despite its simplicity and general efficiency, this strategy alone does not guarantee delivery due to the existence of local minima (or dead ends). Overcoming local minima requires nodes to maintain extra non-local state or to use auxiliary mechanisms. We study how to facilitate greedy forwarding by using a minimum amount of such non-local state in topologically complex networks. Specifically, we investigate the problem of decomposing a given network into a minimum number of Greedily Routable Components (GRC), where greedy routing is guaranteed to work. We approach it by considering an approximate version in a continuous domain, with a central concept called the Greedily Routable Region (GRR). A full characterization of GRR is given concerning its geometric properties and routing capability. We then develop simple approximate algorithms for the problem. These results lead to a practical routing protocol that has a routing stretch below 7 in a continuous domain, and close to 1 in several realistic network settings
Anne-Marie Kermarrec, Guang Tan
MobiHoc1
2010 Connectivity-based and anchor-free localization in large-scale 2d/3d sensor networks
abstract
This paper presents a Connectivity-based and Anchor-free Three-dimensional Localization (CATL) scheme for large-scale sensor networks with concave regions. It distinguishes itself from previous work with a combination of three features: (1) it works for networks in both 2D and 3D spaces, possibly containing holes or concave regions; (2) it is anchor-free, and uses only connectivity information to faithfully recover the original network topology, up to scaling and rotation; (3) it does not depend on the knowledge of network boundaries, which suits it well to situations where boundaries are difficult to identify. The key idea of CATL is to discover the notch nodes, where shortest paths bend and hop-count-based distance starts to significantly deviate from the true Euclidean distance. An iterative protocol is developed that uses a em notch-avoiding multilateration mechanism to localize the network. Simulations show that CATL achieves accurate localization results with a moderate per-node message cost.
Guang Tan, Hongbo Jiang 0001, Shengkai Zhang, Anne-Marie Kermarrec
MobiHoc4
2010 Designing a tit-for-tat based peer-to-peer video-on-demand system
abstract
Video-on-demand (VoD) is a next-generation Internet application of increasing interest allowing users to start watching a movie almost instantaneously by downloading the video on-the-fly. Provided that all users contribute to the system, shifting to the P2P paradigm allows efficient broadcast with a limited-bandwidth source. In VoD applications pieces are downloaded in order. This prevents us from directly applying a BitTorrent-like tit-for-tat incentive scheme. We advocate the use of a loose structure in P2P VoD applications to achieve high playback rates. In this paper we propose a decentralized piece dissemination scheme built on loosely coupled structures maintained using gossip. Peers are grouped into clusters depending on their playback position. Swarming is performed within the clusters while distributed feeding ensures that less advanced clusters get missing pieces from more advanced ones. Our simulations demonstrate that structured dissemination improves from 61% to 77% the achievable playback rate.
Kévin Huguenin, Anne-Marie Kermarrec, Vivek Rai, Maarten van Steen
NOSSDAV2
2010 Application of Random Walks to Decentralized Recommender Systems
Anne-Marie Kermarrec, Vincent Leroy 0001, Afshin Moin, Christopher Thraves
OPODIS1
2010 WhatsUp: News, From, For, Through, Everyone
abstract
WhatsUp (WUP) is a new form of electronic news. It is personalized and decentralized. Users receive news and have the ability to express their interest in it. This opinion, in turn, is used as an implicit and dynamic subscription scheme to filter and personalize future information. The system is peer-to-peer: no big brother company controls the news, and no central server makes it vulnerable to failures, censorship or attacks. At the heart of WUP lies the idea of collaborative filtering applied to the dissemination of news: people who liked the same news in the past might as well like the same news in the future: irrelevant news disappear by themselves. The idea is put to work through Beep: a biased epidemic dissemination (gossip) protocol that delivers news to interested users in a timely manner, despite jamming and churn. Beep is dynamically parameterized on a per- user, per-news, and per-dissemination-hop basis. When compared to a classical epidemic dissemination protocol, Beep has two key characteristics: orientation and amplification. Every user forwards the news of interest to a randomly selected set of users largely constituted by those who have similar interests (orientation). Moreover, the size of this set of users depends on the level of interest in the news itself (amplification).
Antoine Boutet, Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec
Peer-to-Peer Computing4
2010 Boosting Gossip for Live Streaming
abstract
Gossip protocols are considered very effective to disseminate information in a large scale dynamic distributed system. Their inherent simplicity makes them easy to implement and deploy. However, whereas their probabilistic guarantees are often enough to disseminate data in the context of low- bandwidth applications, they typically do not suffice for high-bandwidth content dissemination: missing 1% is unacceptable for live streaming. In this paper, we show how the combination of two simple mechanisms copes with this seemingly inherent deficiency of gossip: (i) codec, an erasure coding scheme, and (ii) claim2, a content- request scheme that leverages gossip duplication to diversify the retransmission sources of missing information. We show how these mechanisms can effectively complement each other in a new gossip protocol, gossip++, which retains the simplicity of deployment of plain gossip. In a realistic setting with an average bandwidth capability (800 kbps) close to the stream rate (680 kbps) and 1% message loss, plain gossip can provide at most 99% of the stream. Using gossip++, on the other hand, all nodes can view a perfectly clear stream.
Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Maxime Monod
Peer-to-Peer Computing3
2010 Brief announcement: byzantine agreement with homonyms
abstract
In this work, we address Byzantine agreement in a message passing system with homonyms, i.e. a system with a number l of authenticated identities that is independent of the total number of processes n, in the presence of t < n Byzantine processes.
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Anne-Marie Kermarrec
SPAA4
2010 Computing in Social Networks
Andrei Giurgiu, Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec
SSS4
2009 Route in Mobile WSN and Get Self-deployment for Free
Kévin Huguenin, Anne-Marie Kermarrec, Eric Fleury
DCOSS2
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
DSN3
2009 Surfing Peer-to-Peer IPTV: Distributed Channel Switching
Anne-Marie Kermarrec, Erwan Le Merrer, Yaning Liu, Gwendal Simon
Euro-Par1
2009 Implementing a Register in a Dynamic Distributed System
abstract
Providing distributed processes with concurrent objects is a fundamental service that has to be offered by any distributed system. The classical shared read/write register is one of the most basic ones. Several protocols have been proposed that build an atomic register on top of an asynchronous message-passing system prone to process crashes. In the same spirit, this paper addresses the implementation of a regular register (a weakened form of an atomic register) in an asynchronous dynamic message-passing system. The aim is here to cope with the net effect of the adversaries that are asynchrony and dynamicity (the fact that processes can enter and leave the system). The paper focuses on the class of dynamic systems the churn rate c of which is constant. It presents two protocols, one applicable to synchronous dynamic message passing systems, the other one to eventually synchronous dynamic systems. Both protocols rely on an appropriate broadcast communication service (similar to a reliable broadcast). Each requires a specific constraint on the churn rate c. Both protocols are first presented in an as intuitive as possible way, and are then proved correct.
Roberto Baldoni, Silvia Bonomi, Anne-Marie Kermarrec, Michel Raynal
ICDCS3
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
ICDCS1
2009 Visibility-Graph-Based Shortest-Path Geographic Routing in Sensor Networks
abstract
We study the problem of shortest-path geographic routing in a static sensor network. Existing algorithms often make routing decisions based on node information in local neighborhoods. However, it is shown by Kuhn et al. that such a design constraint results in a highly undesirable lower bound for routing performance: if a best route has length c, then in the worst case a route produced by any localized algorithm has length Omega(c2), which can be arbitrarily worse than the optimal. We present VIGOR, a visibility-graph-based routing protocol that produces routes of length Theta(c). Our design is based on the construction of a much reduced visibility graph, which guides nodes to find near-optimal paths. The per-node protocol overheads in terms of state information and message transmission depend only on the complexity of the field's large topological features, rather than on the network size. Simulation results show that our protocol dramatically outperforms localized protocols such as GPSR and GOAFR+ in both average and worst cases, with reasonable extra overheads.
Guang Tan, Marin Bertier, Anne-Marie Kermarrec
INFOCOM3
2009 Convex Partition of Sensor Networks and Its Use in Virtual Coordinate Geographic Routing
abstract
Virtual coordinate geographic routing is an appealing geographic routing approach for its ability to work without physical location information. We examine two representative such routing protocols, namely NoGeo and BVR, and show through experiments and theoretical analysis their limitation in adapting to complex field topologies, in particular fields with concave holes. Based on the new insights, we propose a distributed convex partition protocol that divides the field to subareas with convex shapes, using only connectivity information. A new geographic routing protocol, called CONVEX, that builds upon the partitioning protocol is then described. Simulations demonstrate significant performance improvement of the new routing protocol over NoGeo and BVR, in terms of transmission stretch and maintenance overheads.
Guang Tan, Marin Bertier, Anne-Marie Kermarrec
INFOCOM3
2009 Heterogeneous Gossip
Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Boris Koldehofe, Martin Mogensen, Maxime Monod, Vivien Quéma
Middleware3
2009 Decentralized Polling with Respectable Participants
Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod
OPODIS3
2009 Navigating the Web 2.0 with Gossple
Anne-Marie Kermarrec
OPODIS1
2009 Phosphite: Guaranteeing Out of Order Download in P2P Video-on-Demand
abstract
We propose Phosphite, a mechanism to preserve out of order download in peer to peer video-on-demand applications, in the presence of selfish peers. In such applications, peers have a natural trend to download blocks in order to start watching videos as soon as possible. Without specific mechanism to enforce a fair amount of out of order download, the last blocks of the video tend to be lost due to peers leaving soon after having downloaded the last blocks thus forcing peers to rely on the central server for re-introducing those lost blocks. This issue can be solved if peers dedicate a portion of their bandwidth for out of order downloads. Yet, this heavily relies on the goodwill of peers to collaborate. Phosphite is a simple yet efficient approach ensuring that all peers dedicate a part of their bandwidth to out of order download. Phosphite relies on a computational challenge where peers are provided with a combination of the requested blocks and other blocks. This forces peers to download out of order blocks to be able to decode the requested blocks. We evaluate Phosphite and show that it successfully prevents the system from losing blocks, even in the presence of selfish peers, thus offering an appealing alternative to state of the art approaches. With Phosphite, the last blocks remain available (with a probability higher than 0.98), while such result cannot be guaranteed (with a probability lower than 0.5) without enforcement mechanism. Phosphite ensures that a peer to peer download is almost always possible, even in the presence of selfish peers.
Mary-Luc Champel, Anne-Marie Kermarrec, Nicolas Le Scouarnec
Peer-to-Peer Computing2
2009 On Tracking Freeriders in Gossip Protocols
abstract
Peer-to-peer content dissemination applications suffer immensely from freeriders, i.e., nodes that do not provide their fair share. The Tit-for-Tat (TfT) incentives have received much attention as they help make such systems more robust against freeriding. However, these rely on an asymmetric component, namely opportunistic pushes, that let peers receive content without sending anything in return. Opportunistic push constitutes the Achilles' heel of TfT-based protocols as illustrated by the fact that all known attacks against them exploit it. This problem becomes even more serious when used by colluding freeriders. In this paper, we discuss the possibility of using accountability to secure gossip-based dissemination protocols based on asymmetric exchanges. The fact that gossip protocols are dynamic and randomized makes our approach robust against collusion and alleviates the need for cryptography. We present the challenges raised by an auditing approach and give insights into how to build a freerider-tracking protocol for gossip-based content dissemination.
Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod
Peer-to-Peer Computing3
2009 On Gossip and Populations
Marin Bertier, Yann Busnel, Anne-Marie Kermarrec
SIROCCO3
2009 FoG: Fighting the Achilles' Heel of Gossip Protocols with Fountain Codes
Mary-Luc Champel, Anne-Marie Kermarrec, Nicolas Le Scouarnec
SSS2
2009 Challenges in Personalizing and Decentralizing the Web: An Overview of GOSSPLE
Anne-Marie Kermarrec
SSS1
2009 Brief Announcement: Towards Secured Distributed Polling in Social Networks
Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod
DISC3
2009 Editorial
Anne-Marie Kermarrec, Maarten van Steen
Comput. Networks1
2009 Rappel: Exploiting interest and network locality to improve fairness in publish-subscribe systems
Jay A. Patel, Etienne Rivière, Indranil Gupta, Anne-Marie Kermarrec
Comput. Networks4
2009 Slicing Distributed Systems
abstract
Peer-to-peer (P2P) architectures are popular for tasks such as collaborative download, VoIP telephony, and backup. To maximize performance in the face of widely variable storage capacities and bandwidths, such systems typically need to shift work from poor nodes to richer ones. Similar requirements are seen in today's large data centers, where machines may have widely variable configurations, loads, and performance. In this paper, we consider the slicing problem, which involves partitioning the participating nodes into k subsets using a one-dimensional attribute, and updating the partition as the set of nodes and their associated attributes change. The mechanism thus facilitates the development of adaptive systems. We begin by motivating this problem statement and reviewing prior work. Existing algorithms are shown to have problems with convergence, manifesting as inaccurate slice assignments, and to adapt slowly as conditions change. Our protocol, Sliver, has provably rapid convergence, is robust under stress and is simple to implement. We present both theoretical and experimental evaluations of the protocol.
Vincent Gramoli, Ymir Vigfusson, Kenneth P. Birman, Anne-Marie Kermarrec, Robbert van Renesse
IEEE Trans. Computers4
2009 Connectivity-Guaranteed and Obstacle-Adaptive Deployment Schemes for Mobile Sensor Networks
abstract
Mobile sensors can relocate and self-deploy into a network. While focusing on the problems of coverage, existing deployment schemes largely oversimplify the conditions for network connectivity: they either assume that the communication range is large enough for sensors in geometric neighborhoods to obtain location information through local communication, or they assume a dense network that remains connected. In addition, an obstacle-free field or full knowledge of the field layout is often assumed. We present new schemes that are not governed by these assumptions, and thus adapt to a wider range of application scenarios. The schemes are designed to maximize sensing coverage and also guarantee connectivity for a network with arbitrary sensor communication/sensing ranges or node densities, at the cost of a small moving distance. The schemes do not need any knowledge of the field layout, which can be irregular and have obstacles/holes of arbitrary shape. Our first scheme is an enhanced form of the traditional virtual-force-based method, which we term the connectivity-preserved virtual force (CPVF) scheme. We show that the localized communication, which is the very reason for its simplicity, results in poor coverage in certain cases. We then describe a floor-based scheme which overcomes the difficulties of CPVF and, as a result, significantly outperforms it and other state-of-the-art approaches. Throughout the paper our conclusions are corroborated by the results from extensive simulations.
Guang Tan, Stephen A. Jarvis, Anne-Marie Kermarrec
IEEE Trans. Mob. Comput.3
2008 On the Deterministic Tracking of Moving Objects with a Binary Sensor Network
Yann Busnel, Leonardo Querzoni, Roberto Baldoni, Marin Bertier, Anne-Marie Kermarrec
DCOSS5
2008 Connectivity-Guaranteed and Obstacle-Adaptive Deployment Schemes for Mobile Sensor Networks
abstract
Mobile sensors can move and self-deploy into a network. While focusing on the problems of coverage, existing deployment schemes mostly over-simplify the conditions for network connectivity: they either assume that the communication range is large enough for sensors in geometric neighborhoods to obtain each other's locationby local communications, or assume a dense network that remains connected. At the same time, an obstacle-free field or full knowledge of the field layout is often assumed. We present new schemes that are not restricted by these assumptions, and thus adapt to a much wider range of application scenarios. While maximizing sensing coverage, our schemes can achieve connectivity for a network with arbitrary sensor communication/sensing ranges or node densities, at the cost of a small moving distance; the schemes do not need any knowledge of the field layout, which can be irregular and have obstacles/holes of arbitrary shape. Simulations results show that the proposed schemes achieve the targeted properties.
Guang Tan, Stephen A. Jarvis, Anne-Marie Kermarrec
ICDCS3
2008 Distributed churn measurement in arbitrary networks
abstract
We adress the problem of estimating in a fully distributed way the dynamism over a network, called the churn. This BA presents, as far as we know, the first distributed method for monitoring churn in arbitrary networks, subject to arbitrary node departure and arrival patterns.
Vincent Gramoli, Anne-Marie Kermarrec, Erwan Le Merrer
PODC2
2008 A fast distributed slicing algorithm
abstract
No abstract available.
Vincent Gramoli, Ymir Vigfusson, Kenneth P. Birman, Anne-Marie Kermarrec, Robbert van Renesse
PODC4
2008 From anarchy to geometric structuring: the power of virtual coordinates
abstract
This note define self-structuring in a large-scale networked system as the ability of the participating entities to collaboratively impose a geometric structure to the network. This refers to assigning virtual coordinates to participating entities and to dividing the entities in several partitions, in such a way that each entity knows to which partition it belongs.
Anne-Marie Kermarrec, Achour Mostéfaoui, Michel Raynal, Gilles Trédan, Aline Carneiro Viana
PODC1
2008 Reliable Broadcast Tolerating Byzantine Faults in a Message-Bounded Radio Network
Marin Bertier, Anne-Marie Kermarrec, Guang Tan
DISC2
2008 Evaluating the Quality of a Network Topology through Random Walks
Anne-Marie Kermarrec, Erwan Le Merrer, Bruno Sericola, Gilles Trédan
DISC1
2008 SOLIST or How to Look for a Needle in a Haystack? A Lightweight Multi-overlay Structure for Wireless Sensor Networks
abstract
In this paper, we consider sensor database systems. Sensors are attached to objects and queries on the objects are operated at the sensor network level. Although queries to such a system might be extremely complex, ensuring efficiently basic functionalities such as broadcast or anycast without any central element is not trivial. In this paper, we provide a suite of *-cast (anycast, k-cast, broadcast) functionalities in a fully decentralized manner. More specifically, we present the design and evaluation of SOLIST, a multi-layer structure for sensors, largely inspired from structured peer-to-peer systems providing such functionalities. The main goal of SOLIST is to limit the overall energy consumption. A type is associated to each sensor, and the *-cast functionalities are implemented at a type granularity regardless of the number of types and their distribution within the network. A typical use of such a system is sensor-based stock management. We evaluate SOLIST through simulations and show that SOLIST achieves a reasonable trade-off between performance and energy consumption.
Yann Busnel, Marin Bertier, Anne-Marie Kermarrec
WiMob3
2007 A comparison of optimistic approaches to collaborative editing of Wiki pages
abstract
Wikis, a popular tool for sharing knowledge, are basically collaborative editing systems. However, existing Wiki systems offer limited support for co-operative authoring, and they do not scale well, because they are based on a centralised architecture. This paper compares the well-known centralised MediaWiki system with several peer-to-peer approaches to editing of wiki pages: an operational transformation approach (MOT2), a commutativity-oriented approach (WOOTO) and a conflict resolution approach (ACF). We evaluate and compare them, according to a number of qualitative and quantitative metrics.
Claudia-Lavinia Ignat, Gérald Oster, Pascal Molli, Michèle Cart, Jean Ferrié, Anne-Marie Kermarrec, Pierre Sutra, Marc Shapiro 0001, Lamia Benmouffok, Jean-Michel Busca, Rachid Guerraoui
CollaborateCom6
2007 Distributed Slicing in Dynamic Systems
abstract
Peer to peer (P2P) systems are moving from application specific architectures to a generic service oriented design philosophy. This raises interesting problems in connection with providing useful P2P middleware services capable of dealing with resource assignment and management in a large-scale, heterogeneous and unreliable environment. The slicing service, has been proposed to allow for an automatic partitioning of P2P networks into groups (slices) that represent a controllable amount of some resource and that are also relatively homogeneous with respect to that resource. In this paper we propose two gossip-based algorithms to solve the distributed slicing problem. The first algorithm speeds up an existing algorithm sorting a set of uniform random numbers. The second algorithm statistically approximates the rank of nodes in the ordering. The scalability, efficiency and resilience to dynamics of both algorithms rely on their gossip-based models. These algorithms are proved viable theoretically and experimentally.
Antonio Fernández 0001, Vincent Gramoli, Ernesto Jiménez, Anne-Marie Kermarrec, Michel Raynal
ICDCS4
2007 Build One, Get One Free: Leveraging the Coexistence of Multiple P2P Overlay Networks
abstract
Many different P2P overlay networks providing various functionalities, targeting specific applications, have been proposed in the past five years. It is now reasonable to consider that multiple overlays may be deployed over a large set of nodes so that the most appropriate overlay might be chosen depending on the application. A physical peer may then host several instances of logical peers belonging to different overlay networks. In this paper, we show that the coexistence of a structured P2P overlay and an unstructured one may be leveraged so that, by building one, the other is automatically constructed as well. More specifically, we show that the randomness provided by an unstructured gossip-based overlay can be used to build the routing table of a structured P2P overlay and the randomness in the numerical proximity links in the structured networks provides the random peer sampling required by gossip-based unstructured overlays. In this paper, we show that maintaining the leaf set of Pastry and the proximity links of an unstructured overlay is enough to build the complete overlays. Simulation results comparing our approach with both a Pastry-like system and a gossip-based unstructured overlay show that we significantly reduce the overlay maintenance overhead without sacrificing the performance.
Balasubramaneyam Maniymaran, Marin Bertier, Anne-Marie Kermarrec
ICDCS3
2007 VoroNet: A scalable object network based on Voronoi tessellations
abstract
In this paper, we propose the design of VoroNet, an object-based peer to peer overlay network relying on Voronoi tessellations, along with its theoretical analysis and experimental evaluation. VoroNet differs from previous overlay networks in that peers are application objects themselves and get identifiers reflecting the semantics of the application instead of relying on hashing functions. This enables a scalable support for efficient search in large collections of data. In VoroNet, objects are organized in an attribute space according to a Voronoi diagram. VoroNet is inspired from the Kleinberg's small-world model where each peer gets connected to close neighbours and maintains an additional pointer to a long-range neighbour. VoroNet improves upon the original proposal as it deals with general object topologies and therefore copes with skewed data distributions. We show that VoroNet can be built and maintained in a fully decentralized way. The theoretical analysis of the system proves that routing in VoroNet can be achieved in a poly-logarithmic number of hops in the size of the system. The analysis is fully confirmed by our experimental evaluation by simulation.
Olivier Beaumont, Anne-Marie Kermarrec, Loris Marchal, Etienne Rivière
IPDPS2
2007 Peer to Peer Multidimensional Overlays: Approximating Complex Structures
Olivier Beaumont, Anne-Marie Kermarrec, Etienne Rivière
OPODIS2
2007 Small-World Networks: From Theoretical Bounds to Practical Systems
François Bonnet 0001, Anne-Marie Kermarrec, Michel Raynal
OPODIS2
2007 Peer counting and sampling in overlay networks based on random walks
Ayalvadi J. Ganesh, Anne-Marie Kermarrec, Erwan Le Merrer, Laurent Massoulié
Distributed Comput.2
2007 Gossip-based peer sampling
abstract
Gossip-based communication protocols are appealing in large-scale distributed applications such as information dissemination, aggregation, and overlay topology management. This paper factors out a fundamental mechanism at the heart of all these protocols: the peer-sampling service. In short, this service provides every node with peers to gossip with. We promote this service to the level of a first-class abstraction of a large-scale distributed system, similar to a name service being a first-class abstraction of a local-area system. We present a generic framework to implement a peer-sampling service in a decentralized manner by constructing and maintaining dynamic unstructured overlays through gossiping membership information itself. Our framework generalizes existing approaches and makes it easy to discover new ones. We use this framework to empirically explore and compare several implementations of the peer-sampling service. Through extensive simulation experiments we show that---although all protocols provide a good quality uniform random stream of peers to each node locally---traditional theoretical assumptions about the randomness of the unstructured overlays as a whole do not hold in any of the instances. We also show that different design decisions result in severe differences from the point of view of two crucial aspects: load balancing and fault tolerance. Our simulations are validated by means of a wide-area implementation.
Márk Jelasity, Spyros Voulgaris, Rachid Guerraoui, Anne-Marie Kermarrec, Maarten van Steen
ACM Trans. Comput. Syst.4
2006 Topic 15: Peer-to-Peer and Web Computing
Henrique João L. Domingos, Anne-Marie Kermarrec, Pascal Felber, Márk Jelasity
Euro-Par2
2006 Peer sharing behaviour in the eDonkey network, and implications for the design of server-less file sharing systems
abstract
In this paper we present an empirical study of a workload gathered by crawling the eDonkey network --- a dominant peer-to-peer file sharing system --- for over 50 days.We first confirm the presence of some known features, in particular the prevalence of free-riding and the Zipf-like distribution of file popularity. We also analyze the evolution of document popularity.We then provide an in-depth analysis of several clustering properties of such workloads. We measure the geographical clustering of peers offering a given file. We find that most files are offered mostly by peers of a single country, although popular files don't have such a clear home country.We then analyze the overlap between contents offered by different peers. We find that peer contents are highly clustered according to several metrics of interest.We propose to leverage this property by allowing peers to search for content without server support, by querying suitably identified semantic neighbours. We find via trace-driven simulations that this approach is generally effective, and is even more effective for rare files. If we further allow peers to query both their semantic neighbours, and in turn their neighbours' neighbours, we attain hit rates as high as over 55% for neighbour lists of size 20.
Sidath B. Handurukande, Anne-Marie Kermarrec, Fabrice Le Fessant, Laurent Massoulié, Simon Patarin
EuroSys2
2006 Peer to peer size estimation in large and dynamic networks: A comparative study
abstract
As the size of distributed systems keeps growing, the peer to peer communication paradigm has been identified as the key to scalability. Peer to peer overlay networks are characterized by their self-organizing capabilities, resilience to failure and fully decentralized control. In a peer to peer overlay, no entity has a global knowledge of the system. As much as this property is essential to ensure the scalability, monitoring the system under such circumstances is a complex task. Yet, estimating the size of the system is core functionality for many distributed applications to parameter setting or monitoring purposes. In this paper, we propose a comparative study between three algorithms that estimate in a fully decentralized way the size of a peer to peer overlay. Candidate approaches are generally applicable irrespective of the underlying structure of the peer to peer overlay. The paper reports the head to head comparison of estimation system size algorithms. The simulations have been conducted using the same simulation framework and inputs and highlight the differences in cost and accuracy of the estimation between the algorithms both in static and dynamic settings
Erwan Le Merrer, Anne-Marie Kermarrec, Laurent Massoulié
HPDC2
2006 GosSkip, an Efficient, Fault-Tolerant and Self Organizing Overlay Using Gossip-based Construction and Skip-Lists Principles
abstract
This paper presents GosSkip, a self organizing and fully distributed overlay that provides a scalable support to data storage and retrieval in dynamic environments. The structure of GosSkip, while initially possibly chaotic, eventually matches a perfect set of Skip-list-like structures, where no hash is used on data attributes, thus preserving semantic locality and permitting range queries. The use of epidemic-based protocols is the key to scalability, fairness and good behavior of the protocol under churn, while preserving the simplicity of the approach and maintaining O(log(N)) state per peer and O(log(N)) routing costs. In addition, we propose a simple and efficient mechanism to exploit the presence of multiple data items on a single physical node. GosSkip's behavior in both a static and a dynamic scenario is further conveyed by experiments with an actual implementation and real traces of a peer to peer workload
Rachid Guerraoui, Sidath B. Handurukande, Kévin Huguenin, Anne-Marie Kermarrec, Fabrice Le Fessant, Etienne Rivière
Peer-to-Peer Computing4
2006 Ordered Slicing of Very Large-Scale Overlay Networks
abstract
Recently there has been an increasing interest to harness the potential of P2P technology to design and build rich environments where services are provided and multiple applications can be supported in a flexible and dynamic manner. In such a context, resource assignment to services and applications is crucial. Current approaches require significant "manual-mode" operations and/or rely on centralized servers to maintain resource availability. Such approaches are neither scalable nor robust enough. Our contribution towards the solution of this problem is proposing and evaluating a gossip-based protocol to automatically partition the available nodes into "slices", also taking into account specific attributes of the nodes. These slices can be assigned to run services or applications in a fully self-organizing but controlled manner. The main advantages of the proposed protocol are extreme scalability and robustness. We present approximative theoretical models and extensive empirical analysis of the proposed protocol
Márk Jelasity, Anne-Marie Kermarrec
Peer-to-Peer Computing2
2006 Peer counting and sampling in overlay networks: random walk methods
abstract
In this article we address the problem of counting the number of peers in a peer-to-peer system, and more generally of aggregating statistics of individual peers over the whole system. This functionality is useful in many applications, but hard to achieve when each node has only a limited, local knowledge of the whole system. We propose two generic techniques to solve this problem. The Random Tour method is based on the return time of a continuous time random walk to the node originating the query. The Sample and Collide method is based on counting the number of random samples gathered until a target number of redundant samples are obtained. It is inspired by the birthday paradox technique of [6], upon which it improves by achieving a target variance with fewer samples. The latter method relies on a sampling sub-routine which returns randomly chosen peers. Such a sampling algorithm is of independent interest. It can be used, for instance, for neighbour selection by new nodes joining the system. We use a continuous time random walk to obtain such samples. We analyse the complexity and accuracy of the two methods. We illustrate in particular how expansion properties of the overlay affect their performance.
Laurent Massoulié, Erwan Le Merrer, Anne-Marie Kermarrec, Ayalvadi J. Ganesh
PODC3
2006 Efficient and Adaptive Epidemic-Style Protocols for Reliable and Scalable Multicast
abstract
Epidemic-style (gossip-based) techniques have recently emerged as a class of scalable and reliable protocols for peer-to-peer multicast dissemination in large process groups. However, popular implementations of epidemic-style dissemination suffer from two major drawbacks: 1) Network overhead: when deployed on a WAN-wide or VPN-wide scale, they generate a large number of packets that transit across the boundaries of multiple network domains (e.g., LANs, subnets, ASs), causing an overload on core network elements such as bridges, routers, and associated links. 2) Lack of adaptivity: they impose the same load on process group members and the network even under reduced failure rates (viz., packet losses, process failures). In this paper, we describe two protocols to address these problems: 1) a hierarchical gossiping protocol and 2) an adaptive dissemination framework (for multicasts) that allows use of any gossiping primitive within it. These protocols work within a virtual peer-to-peer hierarchy called the leaf box hierarchy. Processes can be allocated in a topologically aware manner to the leaf boxes of this structure, so that protocols 1 and 2 produce low traffic across domain boundaries in the network and induce minimal overhead when there are no failures
Indranil Gupta, Anne-Marie Kermarrec, Ayalvadi J. Ganesh
IEEE Trans. Parallel Distributed Syst.2
2005 Topic 15 - Peer-to-Peer and Web Computing
Anne-Marie Kermarrec, Márk Jelasity, Antony I. T. Rowstron, Henrique João L. Domingos
Euro-Par1
2004 The Peer Sampling Service: Experimental Evaluation of Unstructured Gossip-Based Implementations
Márk Jelasity, Rachid Guerraoui, Anne-Marie Kermarrec, Maarten van Steen
Middleware3
2003 Adaptive Gossip-Based Broadcast
abstract
This 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
DSN5
2003 An Evaluation of Scalable Application-Level Multicast Built Using Peer-To-Peer Overlays
abstract
Structured peer-to-peer overlay networks such as CAN, Chord, Pastry, and Tapestry can be used to implement Internet-scale application-level multicast. There are two general approaches to accomplishing this: tree building and flooding. This paper evaluates these two approaches using two different types of structured overlay: 1) overlays which use a form of generalized hypercube routing, e.g., Chord, Pastry and Tapestry, and 2) overlays which use a numerical distance metric to route through a Cartesian hyperspace, e.g., CAN. Pastry and CAN are chosen as the representatives of each type of overlay. To the best of our knowledge, this paper reports the first head-to-head comparison of CAN-style versus Pastry-style overlay networks, using multicast communication workloads running on an identical simulation infrastructure. The two approaches to multicast are independent of overlay network choice, and we provide a comparison of flooding versus tree-based multicast on both overlays. Results show that the tree-based approach consistently outperforms the flooding approach. Finally, for tree-based multicast, we show that Pastry provides better performance than CAN.
Miguel Castro 0001, Michael B. Jones, Anne-Marie Kermarrec, Antony I. T. Rowstron, Marvin Theimer, Helen J. Wang, Alec Wolman
INFOCOM3
2003 SplitStream: high-bandwidth multicast in cooperative environments
abstract
In tree-based multicast systems, a relatively small number of interior nodes carry the load of forwarding multicast messages. This works well when the interior nodes are highly-available, dedicated infrastructure routers but it poses a problem for application-level multicast in peer-to-peer systems. SplitStream addresses this problem by striping the content across a forest of interior-node-disjoint multicast trees that distributes the forwarding load among all participating peers. For example, it is possible to construct efficient SplitStream forests in which each peer contributes only as much forwarding bandwidth as it receives. Furthermore, with appropriate content encodings, SplitStream is highly robust to failures because a node failure causes the loss of a single stripe on average. We present the design and implementation of SplitStream and show experimental results obtained on an Internet testbed and via large-scale network simulation. The results show that SplitStream distributes the forwarding load among all peers and can accommodate peers with different bandwidth capacities while imposing low overhead for forest construction and maintenance.
Miguel Castro 0001, Peter Druschel, Anne-Marie Kermarrec, Animesh Nandi, Antony I. T. Rowstron, Atul Singh
SOSP3
2003 Network Awareness and Failure Resilience in Self-Organising Overlay Networks
abstract
The growth of peer-to-peer applications on the Internet motivates interest in general purpose overlay networks. The construction of overlays connecting a large population of transient nodes poses several challenges. First, connections in the overlays should reflect the underlying network topology, in order to avoid overloading the network and to allow god application performance. Second, connectivity among active nodes of the overlay should be maintained, even in the presence of high failure rates or when a large proportion of nodes are not active. Finally, the cost of using the overlay should be spread evenly among peer nodes for fairness reasons as well as for the sake of application performance. To preserve scalability, we seek solutions to these issues that can be implemented in a fully decentralized manner and rely on local knowledge from each node. In this paper, we propose an algorithm called the localizer which addresses these three key challenges. The localizer refines the overlay in a way that reflects geographic locality so as to reduce network overload. Simultaneously, it helps to evenly balance the number of neighbors of each node in the overlay, thereby sharing the load evenly as well as improving the resilience to random node failures or disconnections. The proposed algorithm is presented and evaluated in the context of an unstructured peer-to-peer overlay network produced using the Scamp protocol. We provide a theoretical analysis of the various aspects of the algorithm. Simulation results based on a realistic network topology model confirm the analysis and demonstrate the localizer efficiency.
Laurent Massoulié, Anne-Marie Kermarrec, Ayalvadi J. Ganesh
SRDS2
2003 NEEM: Network-Friendly Epidemic Multicast
abstract
Epidemic, 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
SRDS5
2003 HA-PSLS: a highly available parallel single-level store system
abstract
Abstract Parallel single‐level store (PSLS) systems integrate a shared virtual memory and a parallel file system. They provide programmers with a global address space including both memory and file data. PSLS systems implemented in a cluster thus represent a natural support for long‐running parallel applications, combining both the natural shared memory programming model and a large and efficient file system. However, the need to tolerate failures in such a system increases with the size of applications. In this paper we present a highly‐available parallel single level store system (HA‐PSLS), which smoothly integrates a backward error recovery high‐availability mechanism into a PSLS system. Our system is able to tolerate multiple transient failures, a single permanent failure, and power cut failures affecting the whole cluster, without requiring any specialized hardware. For this purpose, HA‐PSLS relies on a high degree of integration (and reusability) of high‐availability and standard features. A prototype integrating our high‐availability support has been implemented and we show some performance results in the paper. Copyright © 2003 John Wiley & Sons, Ltd.
Anne-Marie Kermarrec, Christine Morin
Concurr. Comput. Pract. Exp.1
2003 Peer-to-Peer Membership Management for Gossip-Based Protocols
abstract
Gossip-based protocols for group communication have attractive scalability and reliability properties. The probabilistic gossip schemes studied so far typically assume that each group member has full knowledge of the global membership and chooses gossip targets uniformly at random. The requirement of global knowledge impairs their applicability to very large-scale groups. In this paper, we present SCAMP (Scalable Membership protocol), a novel peer-to-peer membership protocol which operates in a fully decentralized manner and provides each member with a partial view of the group membership. Our protocol is self-organizing in the sense that the size of partial views naturally converges to the value required to support a gossip algorithm reliably. This value is a function of the group size, but is achieved without any node knowing the group size. We propose additional mechanisms to achieve balanced view sizes even with highly unbalanced subscription patterns. We present the design, theoretical analysis, and a detailed evaluation of the basic protocol and its refinements. Simulation results show that the reliability guarantees provided by SCAMP are comparable to previous schemes based on global knowledge. The scale of the experiments attests to the scalability of the protocol.
Ayalvadi J. Ganesh, Anne-Marie Kermarrec, Laurent Massoulié
IEEE Trans. Computers2
2003 Lightweight probabilistic broadcast
abstract
Gossip-based broadcast algorithms, a family of probabilistic broadcast algorithms, trade reliability guarantees against "scalability" properties. Scalability in this context has usually been expressed in terms of message throughput and delivery latency, but there has been little work on how to reduce the memory consumption for membership management and message buffering at large scale.This paper presents lightweight probabilistic broadcast ( lpbcast ), a novel gossip-based broadcast algorithm, which complements the inherent throughput scalability of traditional probabilistic broadcast algorithms with a scalable memory management technique. Our algorithm is completely decentralized and based only on local information: in particular, every process only knows a fixed subset of processes in the system and only buffers fixed "most suitable" subsets of messages. We analyze our broadcast algorithm stochastically and compare the analytical results both with simulations and concrete implementation measurements.
Patrick Eugster, Rachid Guerraoui, Sidath B. Handurukande, Petr Kuznetsov, Anne-Marie Kermarrec
ACM Trans. Comput. Syst.5
2003 Probabilistic Reliable Dissemination in Large-Scale Systems
abstract
The growth of the Internet raises new challenges for the design of distributed systems and applications. In the context of group communication protocols, gossip-based schemes have attracted interest as they are scalable, easy to deploy, and resilient to network and process failures. However, traditional gossip-based protocols have two major drawbacks: 1) they rely on each peer having knowledge of the global membership; and 2) being oblivious to the network topology, they can impose a high load on network links when applied to wide-area settings. In this paper, we provide a theoretical analysis of gossip-based protocols which relates their reliability to key system parameters (the system size, failure rates, and number of gossip targets). The results provide guidelines for the design of practical protocols. In particular, they show how reliability can be maintained while alleviating drawback by: 1) providing each peer with only a small subset of the total membership information and drawback; and 2) organizing members into a hierarchical structure that reflects their proximity according to some network-related metric. We validate the analytical results by simulations and verify that the hierarchical gossip protocol considerably reduces the load on the network compared to the original, non-hierarchical protocol.
Anne-Marie Kermarrec, Laurent Massoulié, Ayalvadi J. Ganesh
IEEE Trans. Parallel Distributed Syst.1
2002 Workshop on Reliable Peer-to-Peer Distributed Systems
Özalp Babaoglu, Anne-Marie Kermarrec, Robbert van Renesse, Luís E. T. Rodrigues, Maarten van Steen, Amin Vadhat
SRDS2
2002 Efficient Epidemic-Style Protocols for Reliable and Scalable Multicast
abstract
Epidemic-style (gossip-based) techniques have recently emerged as a scalable class of protocols for peer-to-peer reliable multicast dissemination in large process groups. These protocols provide probabilistic guarantees on reliability and scalability. However, popular implementations of epidemic-style dissemination are reputed to suffer from two major drawbacks: (a) (Network Overhead) when deployed on a WAN-wide or VPN-wide scale they generate a large number of packets that transit across the boundaries of multiple network domains (e.g., LANs, subnets, ASs), causing an overload on core network elements such as bridges, routers, and associated links; (b) (Lack of Adaptivity) they impose the same load on process group members and the network even under reduced failure rates (viz., packet losses, process failures). lit this paper we report on the (first) comprehensive set of solutions to these problems. The solution is comprised of two protocols: (1) a hierarchical gossiping protocol, and (2) an adaptive multicast dissemination framework that allows use of any gossiping primitive within it. These protocols work within a virtual peer-to-peer hierarchy called the Leaf Box hierarchy. Processes can be allocated in a topologically aware manner to the leaf boxes of this structure, so that (1) and (2) produce low traffic across domain boundaries in the network. In the interests of space, this paper focuses on a detailed discussion and evaluation (through simulations) of only the hierarchical gossiping protocol. We present an overview of the adaptive dissemination protocol and its properties.
Indranil Gupta, Anne-Marie Kermarrec, Ayalvadi J. Ganesh
SRDS2
2002 Scribe: a large-scale and decentralized application-level multicast infrastructure
abstract
This paper presents Scribe, a scalable application-level multicast infrastructure. Scribe supports large numbers of groups, with a potentially large number of members per group. Scribe is built on top of Pastry, a generic peer-to-peer object location and routing substrate overlayed on the Internet, and leverages Pastry's reliability, self-organization, and locality properties. Pastry is used to create and manage groups and to build efficient multicast trees for the dissemination of messages to each group. Scribe provides best-effort reliability guarantees, and we outline how an application can extend Scribe to provide stronger reliability. Simulation results, based on a realistic network topology model, show that Scribe scales across a wide range of groups and group sizes. Also, it balances the load on the nodes while achieving acceptable delay and link stress when compared with Internet protocol multicast.
Miguel Castro 0001, Peter Druschel, Anne-Marie Kermarrec, Antony I. T. Rowstron
IEEE J. Sel. Areas Commun.3
2001 A Two-Level Checkpoint Algorithm in a Highly-Available Parallel Single Level Store System
abstract
A parallel single level store system (PSLS) integrates a shared virtual memory and a parallel file system. Managing the data globally it provides programmers of scientific applications with the attractive shared memory programming model combined with a large and efficient file system in a cluster. We present a cheap and efficient two-level checkpointing approach enabling a PSLS to tolerate failures. The first level checkpointing algorithm is very efficient and saves data in memory but requires a large amount of memory space. When memories are saturated, an alternative algorithm, saving a checkpoint on disks is implemented. Performance results present the impact of different variants of the checkpointing algorithms.
Christine Morin, Renaud Lottiaux, Anne-Marie Kermarrec
CCGRID3
2001 Lightweight Probabilistic Broadcast
abstract
The growing interest in peer-to-peer applications has underlined the importance of scalability in modern distributed systems. Not surprisingly, much research effort has been invested in gossip-based broadcast protocols. These trade the traditional strong reliability guarantees against very good "scalability" properties. Scalability is in that context usually expressed in terms of throughput and delivery latency, but there is only little work on how to reduce the overhead of membership management on a large scale. The paper presents Lightweight Probabilistic Broadcast (lpbcast), a novel gossip-based broadcast algorithm which preserves the inherent throughput scalability of traditional gossip-based algorithms and adds a notion of membership management scalability: every process only knows a random subset of fixed size of the processes in the system. We formally analyze our broadcast algorithm in terms of scalability with respect to the size of individual views, and compare the analytical results both with simulations and concrete measurements.
Patrick Eugster, Rachid Guerraoui, Sidath B. Handurukande, Petr Kuznetsov, Anne-Marie Kermarrec
DSN5
2001 Topic 09: Distributed Systems and Algorithms
Bertil Folliot, Giovanni Chiola, Peter Druschel, Anne-Marie Kermarrec
Euro-Par4
2001 Smooth and Efficient Integration of High-Availability in a Parallel Single Level Store System
Anne-Marie Kermarrec, Christine Morin
Euro-Par1
2001 Probabilistic Semantically Reliable Multicast
abstract
Traditional 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
NCA4
2001 The IceCube approach to the reconciliation of divergent replicas
abstract
We describe a novel approach to log-based reconciliation called IceCube. It is general and is parameterised by application and object semantics. IceCube considers more flexible orderings and is designed to ease the burden of reconciliation on the application programmers. IceCube captures the static and dynamic reconciliation constraints between all pairs of actions, proposes schedules that satisfy the static constraints, and validates them against the dynamic constraints.
Anne-Marie Kermarrec, Antony I. T. Rowstron, Marc Shapiro 0001, Peter Druschel
PODC1
2001 Reducing Noise in Gossip-Based Reliable Broadcast
abstract
We present in this paper a general garbage collection scheme that reduces the "noise" in gossip-based broadcast algorithms. In short, our garbage collection scheme uses a simple heuristic to trade "useless" messages with "useful" ones. Used with a given gossip-based broadcast algorithm, a given size of buffers, and a given number of disseminated messages (e.g., per gossip round), our garbage collection scheme provides higher overall reliability than more conventional schemes. We illustrate our approach through two algorithms: bimodal multicast (pbcast) and lightweight probabilistic broadcast (lpbcast). Our scheme is based on the intuitive idea of discarding messages according to their "age". The "age" of a message represents the number of times the message has been retransmitted.
Petr Kuznetsov, Rachid Guerraoui, Sidath B. Handurukande, Anne-Marie Kermarrec
SRDS4
2000 Improvement of the QoS via an Adaptive and Dynamic Distribution of Applications in a Mobile Environment
abstract
Mobile computing is a domain in great expansion. Wireless networks and Portable Information Appliances (PIAs) are developing very rapidly. More and more mobile users would like to perform their multimedia applications with the same facility as on their desktop station. Use of such applications in a mobile environment raises new challenges. These applications are interactive and extremely costly in system and network resources, whereas PIA resources are poor and wireless networks offer a very variable quality of connection. The authors propose an adaptive and dynamic distribution of applications on the local environment to overcome the poorness of available resources on PIAs, and to reduce and regulate variability effects.
Françoise André, Anne-Marie Kermarrec, Frédéric Le Mouël
SRDS2
2000 High Availability of the Memory Hierarchy in a Cluster
abstract
A single-level store (SLS) integrating a shared virtual memory and a parallel file system with file mapping as its interface is attractive for the execution of high-performance applications in a cluster. However, the probability of a node reboot or failure is quite high. In this paper, we present the design of a highly available SLS system. Our approach combines checkpointing in memory and permanent checkpointing on disk in a cluster using all cluster memory and disk resources. Preliminary performance results show the applicability of the proposed approach for parallel applications with huge input/output requirements.
Christine Morin, Renaud Lottiaux, Anne-Marie Kermarrec
SRDS3
2000 An Efficient and Scalable Approach for Implementing Fault-Tolerant DSM Architectures
abstract
Distributed Shared Memory (DSM) architectures are attractive to execute high performance parallel applications. Made up of a large number of components, these architectures have however a high probability of failure. We propose a protocol to tolerate node failures in cache-based DSM architectures. The proposed solution is based on backward error recovery and consists of an extension to the existing coherence protocol to manage data used by processors for the computation and recovery data used for fault tolerance. This approach can be applied to both Cache Only Memory Architectures (COMA) and Shared Virtual Memory (SVM) systems. The implementation of the protocol in a COMA architecture has been evaluated by simulation. The protocol has also been implemented in an SVM system on a network of workstations. Both simulation results and measurements show that our solution is efficient and scalable.
Christine Morin, Anne-Marie Kermarrec, Michel Banâtre, Alain Gefflaut
IEEE Trans. Computers2
1999 Improving Level of Service for Mobile Users using Context-Awareness
abstract
The development of mobile computing combined, with the exponential growth of the Internet leads to new challenges in the development of large-scale information systems. As they move, users may experience dramatic variations in their environment, in terms of latency, network bandwidth, available services around etc. As a consequence, the level of service of the provided information may strongly depend on the context from which a user issues a query. To handle such variations while providing the best level of service to a user, information systems must be adaptive to change their behavior, preferably without the user intervention, depending on the current context of the user. We propose a general infrastructure based on contextual objects to design and develop adaptive distributed information systems in order to keep, even to improve, the level of the delivered service despite environment variations. As a first step to a complete implementation of our framework, we have implemented a location-aware Web service based on mobile-IP.
Paul Couderc, Anne-Marie Kermarrec
SRDS2
1998 A Framework for Consistent, Replicated Web Objects
abstract
Despite the extensive use of caching techniques, the Web is overloaded. While the caching techniques currently used help some, it would be better to use different caching and replication strategies for different Web pages, depending on their characteristics. We propose a framework in which such strategies can be devised independently per Web document. A Web document is constructed as a worldwide, scalable distributed Web object. Depending on the coherence requirements for that document, the most appropriate caching or replication strategy can subsequently be implemented and encapsulated by the Web object. Coherence requirements are formulated from two different perspectives: that of the Web object, and that of clients using the Web object. We have developed a prototype in Java to demonstrate the feasibility of implementing different strategies for different Web objects.
Anne-Marie Kermarrec, Ihor Kuz, Maarten van Steen, Andrew S. Tanenbaum
ICDCS1
1998 Design, Implementation and Evaluation of ICARE: An Efficient Recoverable DSM
abstract
In the light of the increasing throughput of local area networks, Networks Of Workstations (NOWs) which provide a Distributed Shared Memory (DSM) have become a convenient and cheaper alternative to parallel architectures in the framework of parallel scientific applications. However, the probability that a failure occurs in such a system made up of a large number of components must not be neglected, especially for long-running applications. This paper presents the design, implementation and performance evaluation of ICARE, a page-based recoverable DSM implemented on top of an ATM-based NOW running the CHORUS microkernel. ICARE relies on a Backward Error Recovery (BER) mechanism, and provides a way to combine both efficiency and high-availability. The fact that checkpoints are stored in volatile memory provides a low-cost fault-tolerance mechanism, as well as the opportunity to exploit the symbiotic relationship between the data replication implemented in DSM systems and that needed for fault-tolerance. Furthermore, ICARE efficiently implements transparent process rollback recovery. Performance evaluations show the efficiency of the ICARE prototype that implements the proposed algorithms. © 1998 John Wiley & Sons, Ltd.
Anne-Marie Kermarrec, Christine Morin, Michel Banâtre
Softw. Pract. Exp.1
1996 COMA: An Opportunity for Building Fault-Tolerant Scalable Shared Memory Multiprocessors
abstract
Due to the increasing number of their components, Scalable Shared Memory Multiprocessors (SSMMs) have a very high probability of experiencing failures. Tolerating node failures therefore becomes very important for these architectures particularly if they must be used for long-running computations. In this paper, we show that the class of Cache Only Memory Architectures (COMA) are good candidates for building fault-tolerant SSMMs. A backward error recovery strategy can be implemented without significant hardware modification to previously proposed COMA by exploiting their standard replication mechanisms and extending the coherence protocol to transparently manage recovery data. Evaluation of the proposed fault-tolerant COMA is based on execution driven simulations using some of the Splash applications. We show that, for the simulated architecture, the performance degradation caused by fault-tolerance mechanisms varies from 5% in the best case to 35% in the worst case. The standard memory behavior is only slightly perturbed. Moreover, results also show that the proposed scheme preserves the architecture scalability and that the memory overhead remains low for parallel applications using mostly shared data.
Christine Morin, Alain Gefflaut, Michel Banâtre, Anne-Marie Kermarrec
ISCA4