EDBT 2026 Demo / reviewers in the wild / expert
Gabriel Antoniu
dblp:a/GabrielAntoniu
· DBLP profile ↗
101ranked-venue papers
18as first author
13since 2021 · last 2026
0000-0001-6525-3736ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 79 · 16 first-author · 13 since 2021Databases, data management, data science and information retrieval · 8Artificial intelligence and machine learning · 7Applied, interdisciplinary, general and emerging computing · 7Computer networks · 2Security and privacy · 2Graphics, computer vision, multimedia, augmented reality and games · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Evaluating Federated Learning Beyond Simulation: A Deployment-Aware Methodology
Mathis Valli, Cédric Tedeschi, Gabriel Antoniu, Loïc Cudennec, Alexandru Costan |
CCGrid | 3 |
| 2025 | On the Reproducibility Challenges of Federated Learning: Investigating the Gap Between Simulation, Emulation and Real-World DeploymentsabstractFederated Learning (FL) is an emerging paradigm for decentralized training of Machine Learning models. It has been the subject of a large corpus of research due to its innovative approach to handling sensitive data. A common practice in the FL literature is to run simulations on a single compute node to assess the performance of FL algorithms. While simulation enables fast prototyping and validation of algorithmic concepts, it may face limitations in reproducing the real system's performance in heterogeneous environments such as the Computing Continuum, and particularly on resource-constrained Edge devices. Conversely, emulation on distributed testbeds offers more effective means to accurately reproduce the performance of real-world devices. However, to the best of our knowledge, no prior research has investigated the differences between simulation and emulation in FL experiments. In this paper, we study the complementarity of these approaches and discuss their respective challenges, as a first step towards reproducibility of FL experiments. We illustrate our study with a real-life application used as a baseline: an outdoor air quality forecasting framework with real-world sensors. Our results show that simulation can be used to accurately reproduce model performance metrics, while emulation can effectively reproduce the system performance of real-world experiments. Finally, we present a set of lessons learned on the challenges of FL reproducibility and the selection of experimental infrastructures for FL experiments and applications. Cèdric Prigent, Kate Keahey, Alexandru Costan, Loïc Cudennec, Gabriel Antoniu |
CCGrid | 5 |
| 2025 | Efficient distributed continual learning for steering experiments in real-time
Thomas Bouvier, Bogdan Nicolae, Alexandru Costan, Tekin Bicer, Ian T. Foster, Gabriel Antoniu |
Future Gener. Comput. Syst. | 6 |
| 2024 | Efficient Data-Parallel Continual Learning with Asynchronous Distributed Rehearsal BuffersabstractDeep learning has emerged as a powerful method for extracting valuable information from large volumes of data. However, when new training data arrives continuously (i.e., is not fully available from the beginning), incremental training suffers from catastrophic forgetting (i.e., new patterns are reinforced at the expense of previously acquired knowledge). Training from scratch each time new training data becomes available would result in extremely long training times and massive data accumulation. Rehearsal-based continual learning has shown promise for addressing the catastrophic forgetting challenge, but research to date has not addressed performance and scalability. To fill this gap, we propose an approach based on a distributed rehearsal buffer that efficiently complements data-parallel training on multiple GPUs, allowing us to achieve short runtime and scalability while retaining high accuracy. It leverages a set of buffers (local to each GPU) and uses several asynchronous techniques for updating these local buffers in an embarrassingly parallel fashion, all while handling the communication overheads necessary to augment input mini-batches (groups of training samples fed to the model) using unbiased, global sampling. In this paper we explore the benefits of this approach for classification models. We run extensive experiments on up to 128 GPUs of the ThetaGPU supercomputer to compare our approach with baselines representative of training-from-scratch (the upper bound in terms of accuracy) and incremental training (the lower bound). Results show that rehearsal-based continual learning achieves a top-5 classification accuracy close to the upper bound, while simultaneously exhibiting a runtime close to the lower bound. Thomas Bouvier, Bogdan Nicolae, Hugo Chaugier, Alexandru Costan, Ian T. Foster, Gabriel Antoniu |
CCGrid | 6 |
| 2024 | Simulation of Large-Scale HPC Storage Systems: Challenges and MethodologiesabstractAs the scale of production HPC platforms increases, so does the computing and I/O performance gap, exacerbating the storage bottleneck. High-performance storage systems have been developed to alleviate this bottleneck, but many questions remain concerning their architecture, implementation, and configuration. Answering these questions via experimental campaigns proves arduous. First, some answers are required before deploying the system. Second, once a system hits production the experimental scope is limited by the system's specific configuration and by constraints of production use. In this work we identify challenges posed by the design and validation of a storage simulator. We then propose solutions implemented in Fives, a simulator of HPC workloads on platforms that comprise a parallel file system. We show how our simulator can be instantiated and calibrated for the accurate simulation of a production Lustre deployment. Julien Monniot, Francois Tessier, Henri Casanova, Gabriel Antoniu |
HiPC | 4 |
| 2024 | Efficient Resource-Constrained Federated Learning Clustering with Local Data Compression on the Edge-to-Cloud ContinuumabstractFederated Learning (FL) has been proposed as a privacy-preserving approach for distributed learning over decen-tralized resources. While it can be a highly efficient tool for large-scale collaborative training of Machine Learning (ML) models, its efficiency may be strongly impacted by a high variability in data distributions among clients. Clustered FL tackles this problem by grouping clients with similar data distributions and training personalized models. Despite increasing model accuracy for federated peers, existing clustering approaches overlook system and infrastructure constraints leading to sustainability problems for resource-constrained devices. This paper introduces a new method for resource-constrained FL clustering. We leverage pre-trained autoencoders to compress client data into low dimensional space and build lightweight em-bedding vectors used to cluster federated clients. A randomized quantization approach specifically secures the client embedding vectors against data reconstruction. Extensive experiments using a multi-GPU testbed with multiple scenarios introducing concept drift between clients demonstrate the generalitity of our approach to personalized FL. By minimizing the overall system overhead and improving the model convergence, our approach reduces model training cost by up to l.44× − 4.32× communication, l.03× − 2.40× training time compared to IFCA and l.0× − 8.60× communication, 0.87× − 0.87× training time compared to LADD to achieve similar accuracy in the different evaluation scenarios. While each of the baselines encounters performance degradation in at least one of the scenarios, our strategy demonstrates top efficiency in all of them. Cèdric Prigent, Melvin Chelli, Alexandru Costan, Loïc Cudennec, René Schubotz, Gabriel Antoniu |
HiPC | 6 |
| 2024 | Enabling federated learning across the computing continuum: Systems, challenges and future directions
Cèdric Prigent, Alexandru Costan, Gabriel Antoniu, Loïc Cudennec |
Future Gener. Comput. Syst. | 3 |
| 2023 | FedGuard: Selective Parameter Aggregation for Poisoning Attack Mitigation in Federated LearningabstractMinimizing the attack surface of Federated Learning (FL) systems is a field of active research. FL turns out to be highly vulnerable to various threats coming from the edge of the network. Current approaches rely on robust aggregation, anomaly detection and generative models for defending against poisoning attacks. Yet, they either have limited defensive capabilities due to their underlying design or are impractical to use as they rely on constraining building blocks.We introduce FedGuard, a novel FL framework that utilizes the generative capabilities of Conditional Variational AutoEncoders (CVAE) to effectively defend against poisoning attacks with tuneable overhead in communication and computation. Whilst the idea of hardening a FL system using generative models is not entirely new, FedGuard’s original contribution is in its selective parameter aggregation operator with parameter selection being driven by synthetic validation data sampled from the CVAEs trained locally by each participating party.Experimental evaluations in a 100-client setup demonstrates FedGuard to be more effective than previous approaches against several types of attacks (label and sign flipping, additive noise, same value attacks). FedGuard successfully defends in scenarios with up to 50% malicious peers where other strategies fail. In addition, FedGuard does not require auxiliary datasets or centralized (pre-) training. It provides resilience against poisoning attacks from the very first round of federated training. Melvin Chelli, Cèdric Prigent, René Schubotz, Alexandru Costan, Gabriel Antoniu, Loïc Cudennec, Philipp Slusallek |
CLUSTER | 5 |
| 2023 | ProvLight: Efficient Workflow Provenance Capture on the Edge-to-Cloud ContinuumabstractModern scientific workflows require hybrid infrastructures combining numerous decentralized resources on the IoT/Edge interconnected to Cloud/HPC systems (aka the Computing Continuum) to enable their optimized execution. Understanding and optimizing the performance of such complex Edge-to-Cloud workflows is challenging. Capturing the provenance of key performance indicators, with their related data and processes, may assist in understanding and optimizing workflow executions. However, the capture overhead can be prohibitive, particularly in resource-constrained devices, such as the ones on the IoT/Edge.To address this challenge, based on a performance analysis of existing systems, we propose ProvLight, a tool to enable efficient provenance capture on the IoT/Edge. We leverage simplified data models, data compression and grouping, and lightweight transmission protocols to reduce overheads. We further integrate ProvLight into the E2Clab framework to enable workflow provenance capture across the Edge-to-Cloud Continuum. This integration makes E2Clab a promising platform for the performance optimization of applications through reproducible experiments.We validate ProvLight at a large scale with synthetic workloads on 64 real-life IoT/Edge devices in the FIT IoT LAB testbed. Evaluations show that ProvLight outperforms state-of-the-art systems like ProvLake and DfAnalyzer in resource-constrained devices. ProvLight is 26—37x faster to capture and transmit provenance data; uses 5—7x less CPU; 2x less memory; transmits 2x less data; and consumes 2—2.5x less energy. ProvLight [1] and E2Clab [2] are available as open-source tools. Daniel Rosendo, Marta Mattoso, Alexandru Costan, Renan Souza 0001, Débora B. Pina, Patrick Valduriez, Gabriel Antoniu |
CLUSTER | 7 |
| 2023 | Supporting dynamic allocation of heterogeneous storage resources on HPC systemsabstractSummary Scaling up large‐scale scientific applications on supercomputing facilities is largely dependent on the ability to scale up efficiently data storage and retrieval. However, there is an ever‐widening gap between I/O and computing performance. To address this gap, an increasingly popular approach consists in introducing new intermediate storage tiers (node‐local storage, burst‐buffers,) between the compute nodes and the traditional global shared parallel file‐system. Unfortunately, without advanced techniques to allocate and size these resources, they remain underutilized. In this article, we investigate how heterogeneous storage resources can be allocated on an high‐performance computing platform, just like compute resources. To this purpose, we introduce StorAlloc, a simulator used as a testbed for assessing storage‐aware job scheduling algorithms and evaluating various storage infrastructures. We illustrate its usefulness by showing through a large series of experiments how this tool can be used to size a burst‐buffer partition on a top‐tier supercomputer by using the job history of a production year. Julien Monniot, Francois Tessier, Matthieu Robert, Gabriel Antoniu |
Concurr. Comput. Pract. Exp. | 4 |
| 2022 | Distributed intelligence on the Edge-to-Cloud Continuum: A systematic literature review
Daniel Rosendo, Alexandru Costan, Patrick Valduriez, Gabriel Antoniu |
J. Parallel Distributed Comput. | 4 |
| 2021 | Virtual Log-Structured Storage for High-Performance StreamingabstractOver the past decade, given the higher number of data sources (e.g., Cloud applications, Internet of things) and critical business demands, Big Data transitioned from batch-oriented to real-time analytics. Stream storage systems, such as Apache Kafka, are well known for their increasing role in real-time Big Data analytics. For scalable stream data ingestion and processing, they logically split a data stream topic into multiple partitions. Stream storage systems keep multiple data stream copies to protect against data loss while implementing a stream partition as a replicated log. This architectural choice enables simplified development while trading cluster size with performance and the number of streams optimally managed. This paper introduces a shared virtual log-structured storage approach for improving the cluster throughput when multiple producers and consumers write and consume in parallel data streams. Stream partitions are associated with shared replicated virtual logs transparently to the user, effectively separating the implementation of stream partitioning (and data ordering) from data replication (and durability). We implement the virtual log technique in the KerA stream storage system. When comparing with Apache Kafka, KerA improves the cluster ingestion throughput by up to 4x when multiple producers write over hundreds of data streams. Ovidiu-Cristian Marcu, Alexandru Costan, Bogdan Nicolae, Gabriel Antoniu |
CLUSTER | 4 |
| 2021 | Reproducible Performance Optimization of Complex Applications on the Edge-to-Cloud ContinuumabstractIn more and more application areas, we are witnessing the emergence of complex workflows that combine computing, analytics and learning. They often require a hybrid execution infrastructure with IoT devices interconnected to cloud/HPC systems (aka Computing Continuum). Such workflows are subject to complex constraints and requirements in terms of performance, resource usage, energy consumption and financial costs. This makes it challenging to optimize their configuration and deployment. We propose a methodology to support the optimization of real-life applications on the Edge-to-Cloud Continuum. We implement it as an extension of E2Clab, a previously proposed framework supporting the complete experimental cycle across the Edge-to-Cloud Continuum. Our approach relies on a rigorous analysis of possible configurations in a controlled testbed environment to understand their behaviour and related performance tradeoffs. We illustrate our methodology by optimizing Pl@ntNet, a world-wide plant identification application. Our methodology can be generalized to other applications in the Edge-to-Cloud Continuum. Daniel Rosendo, Alexandru Costan, Gabriel Antoniu, Matthieu Simonin, Jean-Christophe Lombardo, Alexis Joly, Patrick Valduriez |
CLUSTER | 3 |
| 2020 | A Distributed Multi-Sensor Machine Learning Approach to Earthquake Early WarningabstractOur research aims to improve the accuracy of Earthquake Early Warning (EEW) systems by means of machine learning. EEW systems are designed to detect and characterize medium and large earthquakes before their damaging effects reach a certain location. Traditional EEW methods based on seismometers fail to accurately identify large earthquakes due to their sensitivity to the ground motion velocity. The recently introduced high-precision GPS stations, on the other hand, are ineffective to identify medium earthquakes due to its propensity to produce noisy data. In addition, GPS stations and seismometers may be deployed in large numbers across different locations and may produce a significant volume of data consequently, affecting the response time and the robustness of EEW systems.In practice, EEW can be seen as a typical classification problem in the machine learning field: multi-sensor data are given in input, and earthquake severity is the classification result. In this paper, we introduce the Distributed Multi-Sensor Earthquake Early Warning (DMSEEW) system, a novel machine learning-based approach that combines data from both types of sensors (GPS stations and seismometers) to detect medium and large earthquakes. DMSEEW is based on a new stacking ensemble method which has been evaluated on a real-world dataset validated with geoscientists. The system builds on a geographically distributed infrastructure, ensuring an efficient computation in terms of response time and robustness to partial infrastructure failures. Our experiments show that DMSEEW is more accurate than the traditional seismometer-only approach and the combined-sensors (GPS and seismometers) approach that adopts the rule of relative strength. Kevin Fauvel, Daniel Balouek-Thomert, Diego Melgar, Pedro Silva 0007, Anthony Simonet, Gabriel Antoniu, Alexandru Costan, Véronique Masson, Manish Parashar, Ivan Rodero, Alexandre Termier |
AAAI | 6 |
| 2020 | Pufferscale: Rescaling HPC Data Services for High Energy Physics ApplicationsabstractUser-space HPC data services are emerging as an appealing alternative to traditional parallel file systems, because of their ability to be tailored to application needs while eliminating unnecessary overheads incurred by POSIX compliance. The High Energy Physics (HEP) community is progressively turning towards such services to enable high-throughput accesses under heavy concurrency to billions of event data produced by instruments and consumed by subsequent analysis workflows. Such services would benefit from the possibility to be rescaled up and down to adapt to changing workloads, as experimental campaigns progress , in order to optimize resource usage. This paper formalizes rescaling a distributed storage system as a multi objective optimization problem considering three criteria: load balance, data balance, and duration of the rescaling operation. We propose a heuristic for rapidly finding a good approximate solution, while allowing users to weight the criteria as needed. The heuristic is evaluated with Pufferscale, a new rescaling manager for microservice-based distributed storage systems. To validate our approach in a real-world ecosystem, we showcase the use of Pufferscale as a means to enable storage malleability in the HEPnOS storage system for HEP applications. Nathanael Cheriere, Matthieu Dorier, Gabriel Antoniu, Stefan M. Wild, Sven Leyffer, Robert B. Ross |
CCGRID | 3 |
| 2020 | E2Clab: Exploring the Computing Continuum through Repeatable, Replicable and Reproducible Edge-to-Cloud ExperimentsabstractDistributed digital infrastructures for computation and analytics are now evolving towards an interconnected ecosystem allowing complex applications to be executed from IoT Edge devices to the HPC Cloud (aka the Computing Continuum, the Digital Continuum, or the Transcontinuum). Understanding end-to-end performance in such a complex continuum is challenging. This breaks down to reconciling many, typically contradicting application requirements and constraints with low-level infrastructure design choices. One important challenge is to accurately reproduce relevant behaviors of a given application workflow and representative settings of the physical infrastructure underlying this complex continuum. In this paper we introduce a rigorous methodology for such a process and validate it through E2Clab. It is the first platform to support the complete analysis cycle of an application on the Computing Continuum: (i) the configuration of the experimental environment, libraries and frameworks; (ii) the mapping between the application parts and machines on the Edge, Fog and Cloud; (iii) the deployment of the application on the infrastructure; (iv) the automated execution; and (v) the gathering of experiment metrics. We illustrate its usage with a real-life application deployed on the Grid'5000 testbed, showing that our framework allows one to understand and improve performance, by correlating it to the parameter settings, the resource usage and the specifics of the underlying infrastructure. Daniel Rosendo, Pedro Silva 0007, Matthieu Simonin, Alexandru Costan, Gabriel Antoniu |
CLUSTER | 5 |
| 2020 | Mission possible: Unify HPC and Big Data stacks towards application-defined blobs at the storage layer
Pierre Matri, Yevhen Alforov, Álvaro Brandón, María S. Pérez 0001, Alexandru Costan, Gabriel Antoniu, Michael Kuhn 0003, Philip H. Carns, Thomas Ludwig 0002 |
Future Gener. Comput. Syst. | 6 |
| 2020 | How fast can one resize a distributed file system?
Nathanael Cheriere, Matthieu Dorier, Gabriel Antoniu |
J. Parallel Distributed Comput. | 3 |
| 2019 | Investigating Edge vs. Cloud Computing Trade-offs for Stream ProcessingabstractThe recent spectacular rise of the Internet of Things and the associated augmentation of the data deluge motivated the emergence of Edge computing as a means to distribute processing from centralized Clouds towards decentralized processing units close to the data sources. This led to new challenges in ways to distribute processing across Cloud-based, Edge-based or hybrid Cloud/Edge-based infrastructures. In particular, a major question is: how much can one improve (or degrade) the performance of an application by performing computation closer to the data sources rather than in the Cloud? This paper proposes a methodology to understand such performance trade-offs and illustrates it through experimental evaluation with two real-life stream processing use-cases executed on fully-Cloud and hybrid Cloud-Edge testbeds using state-of-the-art processing engines for each environment. We derive a set of take-aways for the community, highlighting the limitations of each environment, the scenarios that could benefit from hybrid Edge-Cloud deployments, what relevant parameters impact performance and how. Pedro Silva 0007, Alexandru Costan, Gabriel Antoniu |
IEEE BigData | 3 |
| 2019 | Is it Worth Relaxing Fault Tolerance to Speed Up Decommission in Distributed Storage Systems?abstractEfficient resource utilization is a major concern for large-scale computer platforms. One method used to lower energy consumption and operational cost is to reduce the amount of idle resources. This can be achieved by using malleability, namely, the possibility for resource managers to dynamically increase or decrease the amount of resources of jobs while they are running. Decommissioning (i.e., removing from the cluster) the idle nodes as soon as possible allows the resource manager to quickly reallocate those nodes to other jobs. Challenges appear when such nodes host part of a distributed storage system. Such storage systems may need to transfer large amounts of data before releasing the nodes, in order to ensure data availability and a certain level of fault tolerance. In this paper, we model and evaluate the performance of the decommission operation when relaxing the level of fault tolerance (i.e., the number of replicas) during this operation. Intuitively, this is expected to reduce the amount of data transfers needed before nodes are released, and thus allow nodes to be returned to the resource manager faster. We quantify theoretically how much time and resources are saved by such a fast decommission strategy compared with a standard decommission that does not temporarily reduce the fault-tolerance level. We establish lower bounds for the duration of the different phases of a fast decommission. We use the lower bounds to estimate when fast decommission would be useful to reduce the usage of core-hours and when not. We implement a prototype for fast decommission and experimentally validate the lower bounds on the duration of the operation and confirm in practice our theoretical findings. Nathanael Cheriere, Matthieu Dorier, Gabriel Antoniu |
CCGRID | 3 |
| 2019 | Efficient Scheduling of Scientific Workflows Using Hot Metadata in a Multisite CloudabstractLarge-scale, data-intensive scientific applications are often expressed as scientific workflows (SWfs). In this paper, we consider the problem of efficient scheduling of a large SWf in a multisite cloud, i.e., a cloud with geo-distributed cloud data centers (sites). The reasons for using multiple cloud sites to run a SWf are that data is already distributed, the necessary resources exceed the limits at a single site, or the monetary cost is lower. In a multisite cloud, metadata management has a critical impact on the efficiency of SWf scheduling as it provides a global view of data location and enables task tracking during execution. Thus, it should be readily available to the system at any given time. While it has been shown that efficient metadata handling plays a key role in performance, little research has targeted this issue in multisite cloud. In this paper, we propose to identify and exploit hot metadata (frequently accessed metadata) for efficient SWf scheduling in a multisite cloud, using a distributed approach. We implemented our approach within a scientific workflow management system, which shows that our approach reduces the execution time of highly parallel jobs up to 64 percent and that of the whole SWfs up to 55 percent. Ji Liu 0003, Luis Pineda-Morales, Esther Pacitti, Alexandru Costan, Patrick Valduriez, Gabriel Antoniu, Marta Mattoso |
IEEE Trans. Knowl. Data Eng. | 6 |
| 2018 | TýrFS: Increasing Small Files Access Performance with Dynamic Metadata ReplicationabstractSmall files are known to pose major performance challenges for file systems. Yet, such workloads are increasingly common in a number of Big Data Analytics workflows or large-scale HPC simulations. These challenges are mainly caused by the common architecture of most state-of-the-art file systems needing one or multiple metadata requests before being able to read from a file. Small input file size causes the overhead of this metadata management to gain relative importance as the size of each file decreases. In this paper we propose a set of techniques leveraging consistent hashing and dynamic metadata replication to significantly reduce this metadata overhead. We implement such techniques inside a new file system named TýrFS, built as a thin layer above the Týr object store. We prove that TýrFS increases small file access performance up to one order of magnitude compared to other state-of-the-art file systems, while only causing a minimal impact on file write throughput. Pierre Matri, María S. Pérez 0001, Alexandru Costan, Gabriel Antoniu |
CCGrid | 4 |
| 2018 | KerA: Scalable Data Ingestion for Stream ProcessingabstractBig Data applications are increasingly moving from batch-oriented execution models to stream-based models that enable them to extract valuable insights close to real-time. To support this model, an essential part of the streaming processing pipeline is data ingestion, i.e., the collection of data from various sources (sensors, NoSQL stores, filesystems, etc.) and their delivery for processing. Data ingestion needs to support high throughput, low latency and must scale to a large number of both data producers and consumers. Since the overall performance of the whole stream processing pipeline is limited by that of the ingestion phase, it is critical to satisfy these performance goals. However, state-of-art data ingestion systems such as Apache Kafka build on static stream partitioning and offset-based record access, trading performance for design simplicity. In this paper we propose KerA, a data ingestion framework that alleviate the limitations of state-of-art thanks to a dynamic partitioning scheme and to lightweight indexing, thereby improving throughput, latency and scalability. Experimental evaluations show that KerA outperforms Kafka up to 4x for ingestion throughput and up to 5x for the overall stream processing throughput. Furthermore, they show that KerA is capable of delivering data fast enough to saturate the big data engine acting as the consumer. Ovidiu-Cristian Marcu, Alexandru Costan, Gabriel Antoniu, María S. Pérez 0001, Bogdan Nicolae, Radu Tudoran, Stefano Bortoli |
ICDCS | 3 |
| 2018 | SLoG: Large-Scale Logging Middleware for HPC and Big Data ConvergenceabstractCloud developers traditionally rely on purpose-specific services to provide the storage model they need for an application. In contrast, HPC developers have a much more limited choice, typically restricted to a centralized parallel file system for persistent storage. Unfortunately, these systems often offer low performance when subject to highly concurrent, conflicting I/O patterns. This makes difficult the implementation of inherently concurrent data structures such as distributed shared logs. Yet, this data structure is key to applications such as computational steering, data collection from physical sensor grids, or discrete event generators. In this paper we tackle this issue. We present SLoG, shared log middleware providing a shared log abstraction over a parallel file system, designed to circumvent the aforementioned limitations. We evaluate SLoG's design on up to 100,000 cores of the Theta supercomputer: the results show high append velocity at scale while also providing substantial benefits for other persistent backend storage systems. Pierre Matri, Philip H. Carns, Robert B. Ross, Alexandru Costan, María S. Pérez 0001, Gabriel Antoniu |
ICDCS | 6 |
| 2018 | Tailwind: Fast and Atomic RDMA-based Replication
Yacine Taleb, Ryan Stutsman, Gabriel Antoniu, Toni Cortes |
USENIX ATC | 3 |
| 2018 | New directions in mobile, hybrid, and heterogeneous clouds for cyberinfrastructures
Jesús Carretero 0001, Francisco Javier García Blas, Gabriel Antoniu, Dana Petcu |
Future Gener. Comput. Syst. | 3 |
| 2018 | Keeping up with storage: Decentralized, write-enabled dynamic geo-replication
Pierre Matri, María S. Pérez 0001, Alexandru Costan, Luc Bougé, Gabriel Antoniu |
Future Gener. Comput. Syst. | 5 |
| 2018 | RM-BDP: Resource management for Big Data platforms
Florin Pop, Radu Prodan, Gabriel Antoniu |
Future Gener. Comput. Syst. | 3 |
| 2017 | How fast can one scale down a distributed file system?abstractFor efficient Big Data processing, efficient resource utilization becomes a major concern as large-scale computing infrastructures such as supercomputers or clouds keep growing in size. Naturally, energy and cost savings can be obtained by reducing idle resources. Malleability, which is the possibility for resource managers to dynamically increase or reduce the resources of jobs, appears as a promising means to progress towards this goal. However, state-of-the-art parallel and distributed file systems have not been designed with malleability in mind. This is mainly due to the supposedly high cost of storage decommission, which is considered to involve expensive data transfers. Nevertheless, as network and storage technologies evolve, old assumptions on potential bottlenecks can be revisited. In this study, we evaluate the viability of malleability as a design principle for a distributed file system. We specifically model the duration of the decommission operation, for which we obtain a theoretical lower bound. Then we consider HDFS as a use case and we show that our model can explain the measured decommission times. The existing decommission mechanism of HDFS is good when the network is the bottleneck, but could be accelerated by up to a factor 3 when the storage is the limiting factor. With the highlights provided by our model, we suggest improvements to speed up decommission in HDFS and we discuss open perspectives for the design of efficient malleable distributed file systems. Nathanael Cheriere, Gabriel Antoniu |
IEEE BigData | 2 |
| 2017 | Towards a unified storage and ingestion architecture for stream processingabstractBig Data applications are rapidly moving from a batch-oriented execution model to a streaming execution model in order to extract value from the data in real-time. However, processing live data alone is often not enough: in many cases, such applications need to combine the live data with previously archived data to increase the quality of the extracted insights. Current streaming-oriented runtimes and middlewares are not flexible enough to deal with this trend, as they address ingestion (collection and pre-processing of data streams) and persistent storage (archival of intermediate results) using separate services. This separation often leads to I/O redundancy (e.g., write data twice to disk or transfer data twice over the network) and interference (e.g., I/O bottlenecks when collecting data streams and writing archival data simultaneously). In this position paper, we argue for a unified ingestion and storage architecture for streaming data that addresses the aforementioned challenge. We identify a set of constraints and benefits for such a unified model, while highlighting the important architectural aspects required to implement it in real life. Based on these aspects, we briefly sketch our plan for future work that develops the position defended in this paper. Ovidiu-Cristian Marcu, Alexandru Costan, Gabriel Antoniu, María S. Pérez 0001, Radu Tudoran, Stefano Bortoli, Bogdan Nicolae |
IEEE BigData | 3 |
| 2017 | Exploring Shared State in Key-Value Store for Window-Based Multi-Pattern Streaming AnalyticsabstractInternational audience Ovidiu-Cristian Marcu, Radu Tudoran, Bogdan Nicolae, Alexandru Costan, Gabriel Antoniu, María S. Pérez 0001 |
CCGrid | 5 |
| 2017 | An Empirical Evaluation of How The Network Impacts The Performance and Energy Efficiency in RAMCloudabstractIn-memory storage systems emerged as a de-facto building block for today's large scale Web architectures and Big Data processing frameworks. Many research and engineering efforts have been dedicated to improve their performance and memory efficiency. More recently, such systems can leverage high-performance networks, e.g., Infiniband. To be able to leverage these systems, it is essential to understand the trade-offs induced by the use of high-performance networks. This paper aims to provide empirical evidence of the impact of client's location on the performance and energy consumption of in-memory storage systems. Through a study carried on RAMCloud, we focus on two settings: 1) clients are collocated within the same network as the storage servers (with Infiniband interconnects), 2) clients access the servers from a remote network, through TCP/IP. We compare and discuss aspects related to scalability and power consumption for these two scenarios which correspond to different deployment models for applications making use of in-memory cloud storage systems. Yacine Taleb, Shadi Ibrahim, Gabriel Antoniu, Toni Cortes |
CCGrid | 3 |
| 2017 | Energy-Driven Straggler Mitigation in MapReduce
Tien-Dat Phan, Shadi Ibrahim, Amelie Chi Zhou, Guillaume Pallez, Gabriel Antoniu |
Euro-Par | 5 |
| 2017 | Characterizing Performance and Energy-Efficiency of the RAMCloud Storage SystemabstractMost large popular web applications, like Facebook and Twitter, have been relying on large amounts of in-memory storage to cache data and offer a low response time. As the main memory capacity of clusters and clouds increases, it becomes possible to keep most of the data in the main memory. This motivates the introduction of in-memory storage systems. While prior work has focused on how to exploit the low-latency of in-memory access at scale, there is very little visibility into the energy-efficiency of in-memory storage systems. Even though it is known that main memory is a fundamental energy bottleneck in computing systems (i.e., DRAM consumes up to 40% of a server's power). In this paper, by the means of experimental evaluation, we have studied the performance and energy-efficiency of RAMCloud - a well-known in-memory storage system. We reveal that although RAMCloud is scalable for read-only applications, it exhibits non-proportional power consumption. We also find that the current replication scheme implemented in RAMCloud limits the performance and results in high energy consumption. Surprisingly, we show that replication can also play a negative role in crash-recovery. Yacine Taleb, Shadi Ibrahim, Gabriel Antoniu, Toni Cortes |
ICDCS | 3 |
| 2017 | Enabling fast failure recovery in shared Hadoop clusters: Towards failure-aware scheduling
Orcun Yildiz, Shadi Ibrahim, Gabriel Antoniu |
Future Gener. Comput. Syst. | 3 |
| 2017 | Failure detector abstractions for MapReduce-based systems
Bunjamin Memishi, María S. Pérez 0001, Gabriel Antoniu |
Inf. Sci. | 3 |
| 2016 | Managing hot metadata for scientific workflows on multisite cloudsabstractLarge-scale scientific applications are often expressed as workflows that help defining data dependencies between their different components. Several such workflows have huge storage and computation requirements, and so they need to be processed in multiple (cloud-federated) datacenters. It has been shown that efficient metadata handling plays a key role in the performance of computing systems. However, most of this evidence concern only single-site, HPC systems to date. In this paper, we present a hybrid decentralized/distributed model for handling hot metadata (frequently accessed metadata) in multisite architectures. We couple our model with a scientific workflow management system (SWfMS) to validate and tune its applicability to different real-life scientific scenarios. We show that efficient management of hot metadata improves the performance of SWfMS, reducing the workflow execution time up to 50% for highly parallel jobs and avoiding unnecessary cold metadata operations. Luis Pineda-Morales, Ji Liu 0003, Alexandru Costan, Esther Pacitti, Gabriel Antoniu, Patrick Valduriez, Marta Mattoso |
IEEE BigData | 5 |
| 2016 | Adaptive Performance-Constrained In Situ Visualization of Atmospheric SimulationsabstractWhile many parallel visualization tools now provide in situ visualization capabilities, the trend has been to feed such tools with large amounts of unprocessed output data and let them render everything at the highest possible resolution. This leads to an increased run time of simulations that still have to complete within a fixed-length job allocation. In this paper, we tackle the challenge of enabling in situ visualization under performance constraints. Our approach shuffles data across processes according to its content and filters out part of it in order to feed a visualization pipeline with only a reorganized subset of the data produced by the simulation. Our framework leverages fast, generic evaluation procedures to score blocks of data, using information theory, statistics, and linear algebra. It monitors its own performance and adapts dynamically to achieve appropriate visual fidelity within predefined performance constraints. Experiments on the Blue Waters supercomputer with the CM1 simulation show that our approach enables a 5x speedup with respect to the initial visualization pipeline and is able to meet performance constraints. Matthieu Dorier, Robert Sisneros, Leonardo Arturo Bautista-Gomez, Tom Peterka, Leigh Orf, Lokman Rahmani, Gabriel Antoniu, Luc Bougé |
CLUSTER | 7 |
| 2016 | Spark Versus Flink: Understanding Performance in Big Data Analytics FrameworksabstractBig Data analytics has recently gained increasing popularity as a tool to process large amounts of data on-demand. Spark and Flink are two Apache-hosted data analytics frameworks that facilitate the development of multi-step data pipelines using directly acyclic graph patterns. Making the most out of these frameworks is challenging because efficient executions strongly rely on complex parameter configurations and on an in-depth understanding of the underlying architectural choices. Although extensive research has been devoted to improving and evaluating the performance of such analytics frameworks, most of them benchmark the platforms against Hadoop, as a baseline, a rather unfair comparison considering the fundamentally different design principles. This paper aims to bring some justice in this respect, by directly evaluating the performance of Spark and Flink. Our goal is to identify and explain the impact of the different architectural choices and the parameter configurations on the perceived end-to-end performance. To this end, we develop a methodology for correlating the parameter settings and the operators execution plan with the resource usage. We use this methodology to dissect the performance of Spark and Flink with several representative batch and iterative workloads on up to 100 nodes. Our key finding is that there none of the two framework outperforms the other for all data types, sizes and job patterns. This paper performs a fine characterization of the cases when each framework is superior, and we highlight how this performance correlates to operators, to resource usage and to the specifics of the internal framework design. Ovidiu-Cristian Marcu, Alexandru Costan, Gabriel Antoniu, María S. Pérez 0001 |
CLUSTER | 3 |
| 2016 | On the Root Causes of Cross-Application I/O Interference in HPC Storage SystemsabstractAs we move toward the exascale era, performance variability in HPC systems remains a challenge. I/O interference, a major cause of this variability, is becoming more important every day with the growing number of concurrent applications that share larger machines. Earlier research efforts on mitigating I/O interference focus on a single potential cause of interference (e.g., the network). Yet the root causes of I/O interference can be diverse. In this work, we conduct an extensive experimental campaign to explore the various root causes of I/O interference in HPC storage systems. We use microbenchmarks on the Grid'5000 testbed to evaluate how the applications' access pattern, the network components, the file system's configuration, and the backend storage devices influence I/O interference. Our studies reveal that in many situations interference is a result of bad flow control in the I/O path, rather than being caused by some single bottleneck in one of its components. We further show that interference-free behavior is not necessarily a sign of optimal performance. To the best of our knowledge, our work provides the first deep insight into the role of each of the potential root causes of interference and their interplay. Our findings can help developers and platform owners improve I/O performance and motivate further research addressing the problem across all components of the I/O stack. Orcun Yildiz, Matthieu Dorier, Shadi Ibrahim, Robert B. Ross, Gabriel Antoniu |
IPDPS | 5 |
| 2016 | Týr: blob storage meets built-in transactionsabstractConcurrent Big Data applications often require high-performance storage, as well as ACID (Atomicity, Consistency, Isolation, Durability) transaction support. Although blobs (binary large objects) are an increasingly popular storage model for such applications, state-of-the-art blob storage systems offer no transaction semantics. This demands users to coordinate data access carefully in order to avoid race conditions, inconsistent writes, overwrites and other problems that cause erratic behavior. We argue there is a gap between existing storage solutions and application requirements, which limits the design of transaction-oriented applications. We introduce Týr, the first blob storage system to provide built-in, multiblob transactions, while retaining sequential consistency and high throughput under heavy access concurrency. Týr offers fine-grained random write access to data and in-place atomic operations. Large-scale experiments with a production application from CERN LHC show Týr throughput outperforming state-of-the-art solutions by more than 75%. Pierre Matri, Alexandru Costan, Gabriel Antoniu, Jesús Montes, María S. Pérez 0001 |
SC | 3 |
| 2016 | TomusBlobs: scalable data-intensive processing on Azure cloudsabstractSummary The emergence of cloud computing has brought the opportunity to use large‐scale compute infrastructures for a broader and broader spectrum of applications and users. As the cloud paradigm gets attractive for the ‘elasticity’ in resource usage and associated costs (the users only pay for resources actually used), cloud applications still suffer from the high latencies and low performance of cloud storage services. As Big Data analysis on clouds becomes more and more relevant in many application areas, enabling high‐throughput massive data processing on cloud data becomes a critical issue, as it impacts the overall application performance. In this paper, we address this challenge at the level of cloud storage. We introduce a concurrency‐optimized data storage system (called TomusBlobs), which federates the virtual disks associated to the Virtual Machines running the application code on the cloud. We demonstrate the performance benefits of our solution for efficient data‐intensive processing by building an optimized prototype MapReduce framework for Microsoft's Azure cloud platform on the basis of TomusBlobs. Finally, we specifically address the limitations of state‐of‐the‐art MapReduce frameworks for reduce‐intensive workloads, by proposing MapIterativeReduce as an extension of the MapReduce model. We validate the aforementioned contributions through large‐scale experiments with synthetic benchmarks and with real‐world applications on the Azure commercial cloud by using resources distributed across multiple data centers; they demonstrate that our solutions bring substantial benefits to data‐intensive applications compared with approaches relying on state‐of‐the‐art cloud object storage. Copyright © 2013 John Wiley & Sons, Ltd. Alexandru Costan, Radu Tudoran, Gabriel Antoniu, Goetz Brasche |
Concurr. Comput. Pract. Exp. | 3 |
| 2016 | Cloud computing for data-driven science and engineeringabstractCloud computing for data-driven science and engineering* During the past decade, data-driven science and engineering have emerged as a key paradigm for performing scientific research, enabling innovations through new kinds of experiments that were earlier impossible.Today's science has access to advanced instruments like next generation genome sequencers, gigapixel survey telescopes, and networks of sensors that monitor cyber-physical systems, and these are generating datasets that are growing exponentially in complexity and data volume.Big Data, across all dimensions of volume, velocity, variety, and veracity, are offering unique opportunities to enable scientific discovery as well as novel challenges to scientific platforms.Such dynamic, distributed, and data-intensive applications hold the solutions to vital scientific and societal problems of the 21st century.In order to achieve breakthrough in new knowledge, there is a need to develop data-driven system models, perform analytics at large scales, manage data from instruments and analyses, and share and visualize the results with scientific peers and the society at large.To this end, cloud computing offers a computing model for running such data-intensive scientific and engineering applications.Clouds have democratized resource access to underserved disciplines, making it possible to perform nontrivial scientific explorations for just a few hundred dollars.Clouds are particularly cost-effective for Big Data applications due to their co-location of elastic compute resources with data, and their use of commodity hardware, which economizes on costs for non-high performance computing (HPC) workloads.Many contemporary Big Data platforms that have emerged from online enterprises such as Google and Twitter are also optimized for such commodity hardware, as found in their own data centers.Of course, there are costs associated with data transfer and storage, in keeping with the pay-as-you-go model, that may not be well suited for applications requiring frequent transfer of and long-term storage of terabytes of data.Likewise, it is valuable to understand how data-intensive or even HPC applications that have been developed for computing grids, at one end, and applications developed for workstation tools like MATLAB and R, at the other end, can be effectively run on clouds.These are some of the practical realities that are worth exploring on the relevance of clouds for data-driven scientific applications.In this special issue, we have compiled a set of articles that discusses new research, development, and deployment efforts in running eScience and eEngineering workloads on cloud infrastructures and platforms.The open solicitation, which followed the 3rd Workshop on Scientific Cloud Computing (ScienceCloud), invited research and case studies on a variety of topics relevant to data-driven scientific computing on clouds: use of cloud-based technologies to address innovative compute and data-driven scientific problems that are not well served by current HPC clusters and grids, programming platforms for elastic and Big Data applications, performance and cost-effective computing on clouds, and gaps in diverse cloud fabrics and service offerings, among others.In all, the special issue received 28 articles, of which six were selected for publication after multiple rounds of reviews and revisions.The special issue starts with two articles that explore the runtime platform support required for executing Big Data science on clouds.In TomusBlobs: Scalable Data-Intensive Processing on Azure Clouds [1], the authors address the limitations of data storage within IaaS clouds such as Amazon S3 and Microsoft Azure BLOBs that are, while co-located in the data center, not present in the virtual machines (VMs) and need to be accessed over the network.Their distributed storage on VMs is optimized for concurrent access and elastic scaling, even *Corrections added on 5 February 2016, after first online publication: references to the paper "Pilot-abstractions for distributed data-intensive cloud applications" have been removed. Yogesh L. Simmhan, Lavanya Ramakrishnan, Gabriel Antoniu, Carole A. Goble |
Concurr. Comput. Pract. Exp. | 3 |
| 2016 | On the energy footprint of I/O management in Exascale HPC systems
Matthieu Dorier, Orcun Yildiz, Shadi Ibrahim, Anne-Cécile Orgerie, Gabriel Antoniu |
Future Gener. Comput. Syst. | 5 |
| 2016 | Governing energy consumption in Hadoop through CPU frequency scaling: An analysis
Shadi Ibrahim, Tien-Dat Phan, Alexandra Carpen-Amarie, Houssem Chihoub, Diana Moise, Gabriel Antoniu |
Future Gener. Comput. Syst. | 6 |
| 2016 | JetStream: Enabling high throughput live event streaming on multi-site clouds
Radu Tudoran, Alexandru Costan, Olivier Nano, Ivo Santos, Hakan Soncu, Gabriel Antoniu |
Future Gener. Comput. Syst. | 6 |
| 2016 | OverFlow: Multi-Site Aware Big Data Management for Scientific Workflows on CloudsabstractThe global deployment of cloud datacenters is enabling large scale scientific workflows to improve performance and deliver fast responses. This unprecedented geographical distribution of the computation is doubled by an increase in the scale of the data handled by such applications, bringing new challenges related to the efficient data management across sites. High throughput, low latencies or cost-related trade-offs are just a few concerns for both cloud providers and users when it comes to handling data across datacenters. Existing solutions are limited to cloud-provided storage, which offers low performance based on rigid cost schemes. In turn, workflow engines need to improvise substitutes, achieving performance at the cost of complex system configurations, maintenance overheads, reduced reliability and reusability. In this paper, we introduce OverFlow, a uniform data management system for scientific workflows running across geographically distributed sites, aiming to reap economic benefits from this geo-diversity. Our solution is environment-aware, as it monitors and models the global cloud infrastructure, offering high and predictable data handling performance for transfer cost and time, within and across sites. OverFlow proposes a set of pluggable services, grouped in a data scientist cloud kit. They provide the applications with the possibility to monitor the underlying infrastructure, to exploit smart data compression, deduplication and geo-replication, to evaluate data management costs, to set a tradeoff between money and time, and optimize the transfer strategy accordingly. The system was validated on the Microsoft Azure cloud across its 6 EU and US datacenters. The experiments were conducted on hundreds of nodes using synthetic benchmarks and real-life bio-informatics applications (A-Brain, BLAST). The results show that our system is able to model accurately the cloud performance and to leverage this for efficient data dissemination, being able to reduce the monetary costs and transfer time by up to three times. Radu Tudoran, Alexandru Costan, Gabriel Antoniu |
IEEE Trans. Cloud Comput. | 3 |
| 2016 | Using Formal Grammars to Predict I/O Behaviors in HPC: The Omnisc'IO ApproachabstractThe increasing gap between the computation performance of post-petascale machines and the performance of their I/O subsystem has motivated many I/O optimizations including prefetching, caching, and scheduling. In order to further improve these techniques, modeling and predicting spatial and temporal I/O patterns of HPC applications as they run has become crucial. In this paper we present Omnisc'IO, an approach that builds a grammar-based model of the I/O behavior of HPC applications and uses it to predict when future I/O operations will occur, and where and how much data will be accessed. To infer grammars, Omnisc'IO is based on StarSequitur, a novel algorithm extending Nevill-Manning's Sequitur algorithm. Omnisc'IO is transparently integrated into the POSIX and MPI I/O stacks and does not require any modification in applications or higher-level I/O libraries. It works without any prior knowledge of the application and converges to accurate predictions of any N future I/O operations within a couple of iterations. Its implementation is efficient in both computation time and memory footprint. Matthieu Dorier, Shadi Ibrahim, Gabriel Antoniu, Robert B. Ross |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2015 | Chronos: Failure-aware scheduling in shared Hadoop clustersabstractHadoop emerged as the de facto state-of-the-art system for MapReduce-based data analytics. The reliability of Hadoop systems depends in part on how well they handle failures. Currently, Hadoop handles machine failures by re-executing all the tasks of the failed machines (i.e., executing recovery tasks). Unfortunately, this elegant solution is entirely entrusted to the core of Hadoop and hidden from Hadoop schedulers. The unawareness of failures therefore may prevent Hadoop schedulers from operating correctly towards meeting their objectives (e.g., fairness, job priority) and can significantly impact the performance of MapReduce applications. This paper presents Chronos, a failure-aware scheduling strategy that enables an early yet smart action for fast failure recovery while still operating within a specific scheduler objective. Upon failure detection, rather than waiting an uncertain amount of time to get resources for recovery tasks, Chronos leverages a lightweight preemption technique to carefully allocate these resources. In addition, Chronos considers data locality when scheduling recovery tasks to further improve the performance. We demonstrate the utility of Chronos by combining it with Fifo and Fair schedulers. The experimental results show that Chronos recovers to a correct scheduling behavior within a couple of seconds only and reduces the job completion times by up to 55% compared to state-of-the-art schedulers. Orcun Yildiz, Shadi Ibrahim, Tran Anh Phuong, Gabriel Antoniu |
IEEE BigData | 4 |
| 2015 | Towards Multi-site Metadata Management for Geographically Distributed Cloud WorkflowsabstractWith their globally distributed datacenters, clouds now provide an opportunity to run complex large-scale applications on dynamically provisioned, networked and federated infrastructures. However, there is a lack of tools supporting data intensive applications across geographically distributed sites. For instance, scientific workflows which handle many small files can easily saturate state-of-the-art distributed filesystems based on centralized metadata servers (e.g. HDFS, PVFS). In this paper, we explore several alternative design strategies to efficiently support the execution of existing workflow engines across multi-site clouds, by reducing the cost of metadata operations. These strategies leverage workflow semantics in a 2-level metadata partitioning hierarchy that combines distribution and replication. The system was validated on the Microsoft Azure cloud across 4 EU and US datacenters. The experiments were conducted on 128 nodes using synthetic benchmarks and real-life applications. We observe as much as 28% gain in execution time for a parallel, geo-distributed real-world application (Montage) and up to 50% for a metadata-intensive synthetic benchmark, compared to a baseline centralized configuration. Luis Pineda-Morales, Alexandru Costan, Gabriel Antoniu |
CLUSTER | 3 |
| 2015 | Exploring Energy-Consistency Trade-Offs in Cassandra Cloud Storage SystemabstractApache Cassandra is an open-source cloud storage system that offers multiple types of operation-level consistency including eventual consistency with multiple levels of guarantees and strong consistency. It is being used by many data-center applications (e.g., Facebook and App Scale). Most existing research efforts have been dedicated to exploring trade-offs such as: consistency vs. Performance, consistency vs. Latency and consistency vs. Monetary cost. In contrast, a little work is focusing on the consistency vs. Energy trade-off. As power bills have become a substantial part of the monetary cost for operating a data-center, this paper aims to provide a clearer understanding of the interplay between consistency and energy consumption. Accordingly, a series of experiments have been conducted to explore the implication of different factors on the energy consumption in Cassandra. Our experiments have revealed a noticeable variation in the energy consumption depending on the consistency level. Furthermore, for a given consistency level, the energy consumption of Cassandra varies with the access pattern and the load exhibited by the application. This further analysis indicates that the uneven distribution of the load amongst different nodes also impacts the energy consumption in Cassandra. Finally, we experimentally compare the impact of four storage configuration and data partitioning policies on the energy consumption in Cassandra: interestingly, we achieve 23% energy saving when assigning 50% of the nodes to the hot pool for the applications with moderate ratio of reads and writes, while applying eventual (quorum) consistency. This study points to opportunities for future research on consistency-energy trade-offs and offers useful insight into designing energy-efficient techniques for cloud storage systems. Houssem Chihoub, Shadi Ibrahim, Gabriel Antoniu, María S. Pérez 0001, Luc Bougé |
SBAC-PAD | 4 |
| 2015 | A formal method for rule analysis and validation in distributed data aggregation service
Vlad Serbanescu 0001, Florin Pop, Valentin Cristea, Gabriel Antoniu |
World Wide Web | 4 |
| 2014 | Architecture of Distributed Data Aggregation ServiceabstractThe ever-growing trend of deploying applications over the Internet has resulted in increasingly tougher constraints and requirements. Data management systems are a major concern when it comes to scalability, flexibility and reliability due to being implemented in a distributed way. In this paper we present a Distributed Data Aggregation Service relying on a storage system designed to meet these demands, namely Blob Seer. The primary goal is to serve as a repository backend for complex analysis and automatic mining of scientific data (like bibtex entries). Several requirements, derived from this objective, match Blob Seer's features: versioning used for lock-free access to data and different granularity of read / write operations. We proposed a model to perform the correct translation between Blob Seer's unstructured view of data and the user's structured view. We implemented a client providing a formal description for the data retrieval queries and a specification for a search API. A benchmark tool relying on a performance model of Blob Seer, will be used to automatically determine the best Blob Seer deployment configuration for a specific data aggregation pattern. Vlad Serbanescu 0001, Florin Pop, Valentin Cristea, Gabriel Antoniu |
AINA | 4 |
| 2014 | Bridging Data in the Clouds: An Environment-Aware System for Geographically Distributed Data TransfersabstractInternational audience Radu Tudoran, Alexandru Costan, Luc Bougé, Gabriel Antoniu |
CCGRID | 5 |
| 2014 | Evaluating Streaming Strategies for Event Processing Across Infrastructure CloudsabstractInfrastructure clouds revolutionized the way in which we approach resource procurement by providing an easy way to lease compute and storage resources on short notice, for a short amount of time, and on a pay-as-you-go basis. This new opportunity, however, introduces new performance trade-offs. Making the right choices in leveraging different types of storage available in the cloud is particularly important for applications that depend on managing large amounts of data within and across clouds. An increasing number of such applications conform to a pattern in which data processing relies on streaming the data to a compute platform where a set of similar operations is repeatedly applied to independent chunks of data. This pattern is evident in virtual observatories such as the Ocean Observatory Initiative, in cases when new data is evaluated against existing features in geospatial computations or when experimental data is processed as a series of time events. In this paper, we propose two strategies for efficiently implementing such streaming in the cloud and evaluate them in the context of an ATLAS application processing experimental data. Our results show that choosing the right cloud configuration can improve overall application performance by as much as three times. Radu Tudoran, Kate Keahey, Pierre Riteau, Sergey Panitkin, Gabriel Antoniu |
CCGRID | 5 |
| 2014 | CALCioM: Mitigating I/O Interference in HPC Systems through Cross-Application CoordinationabstractUnmatched computation and storage performance in new HPC systems have led to a plethora of I/O optimizations ranging from application-side collective I/O to network and disk-level request scheduling on the file system side. As we deal with ever larger machines, the interference produced by multiple applications accessing a shared parallel file system in a concurrent manner becomes a major problem. Interference often breaks single-application I/O optimizations, dramatically degrading application I/O performance and, as a result, lowering machine wide efficiency. This paper focuses on CALCioM, a framework that aims to mitigate I/O interference through the dynamic selection of appropriate scheduling policies. CALCioM allows several applications running on a supercomputer to communicate and coordinate their I/O strategy in order to avoid interfering with one another. In this work, we examine four I/O strategies that can be accommodated in this framework: serializing, interrupting, interfering and coordinating. Experiments on Argonne's BG/P Surveyor machine and on several clusters of the French Grid'5000 show how CALCioM can be used to efficiently and transparently improve the scheduling strategy between two otherwise interfering applications, given specified metrics of machine wide efficiency. Matthieu Dorier, Gabriel Antoniu, Robert B. Ross, Dries Kimpe, Shadi Ibrahim |
IPDPS | 2 |
| 2014 | Omnisc'IO: A Grammar-Based Approach to Spatial and Temporal I/O Patterns PredictionabstractThe increasing gap between the computation performance of post-petascale machines and the performance of their I/O subsystem has motivated many I/O optimizations including prefetching, caching, and scheduling techniques. In order to further improve these techniques, modeling and predicting spatial and temporal I/O patterns of HPC applications as they run has became crucial. In this paper we present Omnisc'IO, an approach that builds a grammar-based model of the I/O behavior of HPC applications and uses it to predict when future I/O operations will occur, and where and how much data will be accessed. Omnisc'IO is transparently integrated into the POSIX and MPI I/O stacks and does not require any modification in applications or higher level I/O libraries. It works without any prior knowledge of the application and converges to accurate predictions within a couple of iterations only. Its implementation is efficient in both computation time and memory footprint. Matthieu Dorier, Shadi Ibrahim, Gabriel Antoniu, Robert B. Ross |
SC | 3 |
| 2014 | Transfer as a Service: Towards a Cost-Effective Model for Multi-site Cloud Data ManagementabstractThe global deployment of cloud datacenters is enabling large web services to deliver fast response to users worldwide. This unprecedented geographical distribution of the computation also brings new challenges related to the efficient data management across sites. High throughput, low latencies, cost-or energy-related trade-offs are just a few concerns for both cloud providers and users when it comes to handling data across datacenters. Existing cloud data management solutions are limited to cloud-provided storage, which offers low performance based on rigid cost schemas. In this paper, we are proposing a dedicated cloud data transfer service that supports large-scale data dissemination across geographically distributed sites, advocating for a Transfer as a Service (TaaS) paradigm. The system aggregates the available bandwidth by enabling multiroute transfers across cloud sites. For users of multi-site or federated clouds, our proposal is able to decrease the variability of transfers and increase the throughput up to three times compared to baseline user options, while benefiting from the well-known high availability of cloud-provided services. For cloud providers, such a service can decrease the energy consumption within a datacenter down to half compared to user-based transfers. Radu Tudoran, Alexandru Costan, Gabriel Antoniu |
SRDS | 3 |
| 2013 | Adaptive file management for scientific workflows on the Azure cloudabstractScientific workflows typically communicate data between tasks using files. Currently, on public clouds, this is achieved by using the cloud storage services, which are unable to exploit the workflow semantics and are subject to low throughput and high latencies. To overcome these limitations, we propose an alternative leveraging data locality through direct file transfers between the compute nodes. We rely on the observation that workflows generate a set of common data access patterns that our solution exploits in conjunction with context information to self-adapt, choose the most adequate transfer protocol and expose the data layout within the virtual machines to the workflow engines. This file management system was integrated within the Microsoft Generic Worker workflow engine and was validated using synthetic benchmarks and a real-life application on the Azure cloud. The results show it can bring significant performance gains: up to 5x file transfer speedup compared to solutions based on standard cloud storage and over 25% application timespan reduction compared to Hadoop on Azure. Radu Tudoran, Alexandru Costan, Ramin Rezai Rad, Goetz Brasche, Gabriel Antoniu |
IEEE BigData | 5 |
| 2013 | Consistency in the Cloud: When Money Does Matter!abstractWith the emergence of cloud computing, many organizations have moved their data to the cloud in order to provide scalable, reliable and highly available services. To meet the ever-growing user needs, these services mainly rely on geographically-distributed data replication to guarantee good performance and high availability. However, with replication, consistency comes into question. Service providers in the cloud have the freedom to select the level of consistency according to the access patterns exhibited by the applications. Most optimizations efforts then concentrate on how to provide adequate trade-offs between consistency guarantees and performance. However, as the monetary cost completely relies on the service providers, in this paper we argue that monetary cost should be taken into consideration when evaluating or selecting a consistency level in the cloud. Accordingly, we define a new metric called consistency-cost efficiency. Based on this metric, we present a simple, yet efficient economical consistency model, called Bismar, that adaptively tunes the consistency level at runtime in order to reduce the monetary cost while simultaneously maintaining a low fraction of stale reads. Experimental evaluations with the Cassandra cloud storage on the Grid'5000 test bed show the validity of the metric and demonstrate the effectiveness of the proposed consistency model. Houssem Chihoub, Shadi Ibrahim, Gabriel Antoniu, María S. Pérez 0001 |
CCGRID | 3 |
| 2013 | Distributed Data Storage in Support for Context-Aware ApplicationsabstractContext-aware computing is a new paradigm that relies on large amounts of data collected from a variety of sources, ranging from smartphones to sensors, to automatically take smart decisions. This usually leads to large volumes of data, that need to be further processed to derive higher-level context information. Clouds have recently emerged as interesting candidates to support the storage and aggregation of such data for large-scale context-aware applications. However, specific extensions to support context-aware data need to be designed in order to be able to fully exploit the clouds' potential. In this paper we introduce such a cloud-based system, designed to support real-time processing and persistent storage of context data. Context Aware Framework is designed as an extension of the BlobSeer storage system, building a context-aware layer on top of it to enable scalable high-throughput under high-concurrency for big context data. Our experimental evaluation validates the transparency, mobility and real-time guarantees provided by our approach to context-aware applications. Elena Burceanu, Ciprian Dobre, Valentin Cristea, Alexandru Costan, Gabriel Antoniu |
ISPDC | 5 |
| 2013 | Flubber: Two-level disk scheduling in virtualized environment
Hai Jin 0001, Shadi Ibrahim, Wenzhi Cao, Song Wu 0001, Gabriel Antoniu |
Future Gener. Comput. Syst. | 6 |
| 2013 | GMonE: A complete approach to cloud monitoring
Jesús Montes, Alberto Sánchez 0001, Bunjamin Memishi, María S. Pérez 0001, Gabriel Antoniu |
Future Gener. Comput. Syst. | 5 |
| 2013 | Handling partitioning skew in MapReduce using LEEN
Shadi Ibrahim, Hai Jin 0001, Lu Lu 0006, Bingsheng He, Gabriel Antoniu, Song Wu 0001 |
Peer-to-Peer Netw. Appl. | 5 |
| 2012 | Maestro: Replica-Aware Map Scheduling for MapReduceabstractMapReduce has emerged as a leading programming model for data-intensive computing. Many recent research efforts have focused on improving the performance of the distributed frameworks supporting this model. Many optimizations are network-oriented and most of them mainly address the data shuffling stage of MapReduce. Our studies with Hadoop demonstrate that, apart from the shuffling phase, another source of excessive network traffic is the high number of map task executions which process remote data. That leads to an excessive number of useless speculative executions of map tasks and to an unbalanced execution of map tasks across different machines. All these factors produce a noticeable performance degradation. We propose a novel scheduling algorithm for map tasks, named Maestro, to improve the overall performance of the MapReduce computation. Maestro schedules the map tasks in two waves: first, it fills the empty slots of each data node based on the number of hosted map tasks and on the replication scheme for their input data, second, runtime scheduling takes into account the probability of scheduling a map task on a given machine depending on the replicas of the task's input data. These two waves lead to a higher locality in the execution of map tasks and to a more balanced intermediate data distribution for the shuffling phase. In our experiments on a 100-node cluster, Maestro achieves around 95% local map executions, reduces speculative map tasks by 80% and results in an improvement of up to 34% in the execution time. Shadi Ibrahim, Hai Jin 0001, Lu Lu 0006, Bingsheng He, Gabriel Antoniu, Song Wu 0001 |
CCGRID | 5 |
| 2012 | TomusBlobs: Towards Communication-Efficient Storage for MapReduce Applications in AzureabstractThe emergence of cloud computing brought the opportunity to use large-scale compute infrastructures for a broad spectrum of applications and users. As the cloud paradigm gets attractive for the " elasticity'' in resource usage and associated costs (the users only pay for resources actually used), cloud applications still suffer from the high latencies and low performance of cloud storage services. Enabling high-throughput massive data processing on cloud data becomes a critical issue, as it impacts the overall application performance. In this paper we address the above challenge at the level of the cloud storage. We introduce a concurrency-optimized data storage system which federates the virtual disks associated to VMs. We demonstrate the performance of our solution for efficient data-intensive processing on commercial clouds by building an optimized prototype MapReduce framework for Azure that leverages the benefits of our storage solution. We perform extensive synthetic benchmarks as well as experiments with real-world applications: they demonstrate that our solution brings substantial benefits to data intensive applications compared to approaches relying on state-of-the-art cloud object storage. Radu Tudoran, Alexandru Costan, Gabriel Antoniu, Hakan Soncu |
CCGRID | 3 |
| 2012 | Harmony: Towards Automated Self-Adaptive Consistency in Cloud StorageabstractIn just a few years cloud computing has become a very popular paradigm and a business success story, with storage being one of the key features. To achieve high data availability, cloud storage services rely on replication. In this context, one major challenge is data consistency. In contrast to traditional approaches that are mostly based on strong consistency, many cloud storage services opt for weaker consistency models in order to achieve better availability and performance. This comes at the cost of a high probability of stale data being read, as the replicas involved in the reads may not always have the most recent write. In this paper, we propose a novel approach, named Harmony, which adaptively tunes the consistency level at run-time according to the application requirements. The key idea behind Harmony is an intelligent estimation model of stale reads, allowing to elastically scale up or down the number of replicas involved in read operations to maintain a low (possibly zero) tolerable fraction of stale reads. As a result, Harmony can meet the desired consistency of the applications while achieving good performance. We have implemented Harmony and performed extensive evaluations with the Cassandra cloud storage on Grid'5000 test bed and on Amazon EC2. The results show that Harmony can achieve good performance without exceeding the tolerated number of stale reads. For instance, in contrast to the static eventual consistency used in Cassandra, Harmony reduces the stale data being read by almost 80% while adding only minimal latency. Meanwhile, it improves the throughput of the system by 45% while maintaining the desired consistency requirements of the applications when compared to the strong consistency model in Cassandra. Houssem Chihoub, Shadi Ibrahim, Gabriel Antoniu, María S. Pérez 0001 |
CLUSTER | 3 |
| 2012 | Damaris: How to Efficiently Leverage Multicore Parallelism to Achieve Scalable, Jitter-free I/OabstractWith exascale computing on the horizon, the performance variability of I/O systems represents a key challenge in sustaining high performance. In many HPC applications, I/O is concurrently performed by all processes, which leads to I/O bursts. This causes resource contention and substantial variability of I/O performance, which significantly impacts the overall application performance and, most importantly, its predictability over time. In this paper, we propose a new approach to I/O, called Damaris, which leverages dedicated I/O cores on each multicore SMP node, along with the use of shared-memory, to efficiently perform asynchronous data processing and I/O in order to hide this variability. We evaluate our approach on three different platforms including the Kraken Cray XT5 supercomputer (ranked 11th in Top500), with the CM1 atmospheric model, one of the target HPC applications for the Blue Waters postpetascale supercomputer project. By overlapping I/O with computation and by gathering data into large files while avoiding synchronization between cores, our solution brings several benefits: 1) it fully hides jitter as well as all I/O-related costs, which makes simulation performance predictable, 2) it increases the sustained write throughput by a factor of 15 compared to standard approaches, 3) it allows almost perfect scalability of the simulation up to over 9,000 cores, as opposed to state-of-the-art approaches which fail to scale, 4) it enables a 600% compression ratio without any additional overhead, leading to a major reduction of storage requirements. Matthieu Dorier, Gabriel Antoniu, Franck Cappello, Marc Snir, Leigh Orf |
CLUSTER | 2 |
| 2012 | On-the-Fly Task Execution for Speeding Up Pipelined MapReduce
Diana Moise, Gabriel Antoniu, Luc Bougé |
Euro-Par | 2 |
| 2011 | Managing Data Access on Clouds: A Generic Framework for Enforcing Security PoliciesabstractProviding an adequate security level in Cloud Environments is currently an extremely active research area. More specifically, malicious behaviors targeting large-scale Cloud data repositories (e.g. Denial of Service attacks) may drastically degrade the overall performance of such systems and cannot be detected by typical authentication mechanisms. In this paper we propose a generic security management framework allowing providers of Cloud data management systems to define and enforce complex security policies. This security framework is designed to detect and stop a large array of attacks defined through an expressive policy description language and to be easily interfaced with various data management systems. We show that we can efficiently protect a data storage system by evaluating our security framework on top of the BlobSeer data management platform. We evaluate the benefits of preventing a DoS attack targeted towards BlobSeer through experiments performed on the Grid'5000 testbed. Cristina Basescu, Alexandra Carpen-Amarie, Catalin Adrian Leordeanu, Alexandru Costan, Gabriel Antoniu |
AINA | 5 |
| 2011 | Efficient Support for MPI-I/O Atomicity Based on VersioningabstractWe consider the challenge of building data management systems that meet an important requirement of today's data-intensive HPC applications: to provide a high I/O throughput while supporting highly concurrent data accesses. In this context, many applications rely on MPI-I/O and require atomic, non-contiguous I/O operations that concurrently access shared data. In most existing implementations, the atomicity requirement is implemented through locking-based schemes, which have proven inefficient, especially for non-contiguous I/O. We claim that using a versioning-enabled storage back-end has the potential to avoid the expensive synchronization induced by locking-based schemes. We describe a prototype implementation on top of ROMIO, and report on promising experimental results with standard MPI-I/O benchmarks specifically designed to evaluate the performance of non-contiguous, overlapped I/O accesses under MPI atomicity guarantees. Viet-Trung Tran, Bogdan Nicolae, Gabriel Antoniu, Luc Bougé |
CCGRID | 3 |
| 2011 | Optimizing Multi-deployment on Clouds by Means of Self-adaptive Prefetching
Bogdan Nicolae, Franck Cappello, Gabriel Antoniu |
Euro-Par (1) | 3 |
| 2011 | Introduction
Salvatore Orlando 0001, Gabriel Antoniu, Amol Ghoting, María S. Pérez 0001 |
Euro-Par (1) | 2 |
| 2011 | Going back and forth: efficient multideployment and multisnapshotting on cloudsabstractInfrastructure as a Service (IaaS) cloud computing has revolutionized the way we think of acquiring resources by introducing a simple change: allowing users to lease computational resources from the cloud provider's datacenter for a short time by deploying virtual machines (VMs) on these resources. This new model raises new challenges in the design and development of IaaS middleware. One of those challenges is the need to deploy a large number (hundreds or even thousands) of VM instances simultaneously. Once the VM instances are deployed, another challenge is to simultaneously take a snapshot of many images and transfer them to persistent storage to support management tasks, such as suspend-resume and migration. With datacenters growing rapidly and configurations becoming heterogeneous, it is important to enable efficient concurrent deployment and snapshotting that are at the same time hypervisor independent and ensure a maximum compatibility with different configurations. This paper addresses these challenges by proposing a virtual file system specifically optimized for virtual machine image storage. It is based on a lazy transfer scheme coupled with object versioning that handles snapshotting transparently in a hypervisor-independent fashion, ensuring high portability for different configurations. Large-scale experiments on hundreds of nodes demonstrate excellent performance results: speedup for concurrent VM deployments ranges from a factor of 2 up to 25, with a reduction in bandwidth utilization of as much as 90%. Bogdan Nicolae, John Bresnahan, Kate Keahey, Gabriel Antoniu |
HPDC | 4 |
| 2011 | BlobSeer: Next-generation data management for large scale infrastructures
Bogdan Nicolae, Gabriel Antoniu, Luc Bougé, Diana Moise, Alexandra Carpen-Amarie |
J. Parallel Distributed Comput. | 2 |
| 2010 | Autonomic Cloud Storage: Challenges at StakeabstractWhile the cloud computing paradigm is progressively being adopted by companies aiming to deliver large-scale distributed services, such as Amazon, IBM, Google or Yahoo!, the service level provided for data storage remains rather basic. This talk will discuss several issues related to building advanced facilities for data sharing on cloud infrastructures. We will discuss how open issues raised by the need for efficient, secure and reliable storage service for data intensive distributed applications running in cloud environments may be addressed by enabling an autonomic behavior for the cloud storage infrastructure. Gabriel Antoniu |
CISIS | 1 |
| 2010 | Bringing Introspection Into the BlobSeer Data-Management System Using the MonALISA Distributed Monitoring FrameworkabstractIntrospection is the prerequisite of an autonomic behavior, the first step towards a performance improvement and a resource-usage optimization for large-scale distributed systems. In grid environments, the task of observing the application behavior is assigned to monitoring systems. However, most of them are designed to provide general resource information and do not consider specific information for higher-level services. More specifically, in the context of data-intensive applications, a specific introspection layer is required in order to collect data about the usage of storage resources, about data access patterns, etc. This paper discusses the requirements for an introspection layer in a data-management system for large-scale distributed infrastructures. We focus on the case of BlobSeer, a large-scale distributed system for storing massive data. The paper explains why and how to enhance BlobSeer with introspective capabilities and proposes a three-layered architecture relying on the MonALISA monitoring framework. This approach has been evaluated on the Grid'5000 testbed, with experiments that prove the feasibility of generating relevant information related to the state and the behavior of the system. Alexandra Carpen-Amarie, Alexandru Costan, Gabriel Antoniu, Luc Bougé |
CISIS | 4 |
| 2010 | Using Global Behavior Modeling to Improve QoS in Cloud Data Storage ServicesabstractThe cloud computing model aims to make large-scale data-intensive computing affordable even for users with limited financial resources, that cannot invest into expensive infrastructures necesssary to run them. In this context, MapReduce is emerging as a highly scalable programming paradigm that enables high-throughput data-intensive processing as a cloud service. Its performance is highly dependent on the underlying storage service, responsible to efficiently support massively parallel data accesses by guaranteeing a high throughput under heavy access concurrency. In this context, quality of service plays a crucial role: the storage service needs to sustain a stable throughput for each individual accesss, in addition to achieving a high aggregated throughput under concurrency. In this paper we propose a technique to address this problem using component monitoring, application-side feedback and behavior pattern analysis to automatically infer useful knowledge about the causes of poor quality of service and provide an easy way to reason in about potential improvements. We apply our proposal to Blob Seer, a representative data storage service specifically designed to achieve high aggregated throughputs and show through extensive experimentation substantial improvements in the stability of individual data read accesses under MapReduce workloads. Jesús Montes, Bogdan Nicolae, Gabriel Antoniu, Alberto Sánchez 0001, María S. Pérez 0001 |
CloudCom | 3 |
| 2010 | Improving the Hadoop map/reduce framework to support concurrent appends through the BlobSeer BLOB management systemabstractHadoop is a reference software framework supporting the Map/Reduce programming model. It relies on the Hadoop Distributed File System (HDFS) as its primary storage system. Although HDFS does not offer support for concurrently appending data to existing files, we argue that Map/Reduce applications as well as other classes of applications can benefit from such a functionality. We provide support for concurrent appends by building a concurrency-optimized data storage layer based on the BlobSeer data management service. Moreover, we modify the Hadoop Map/Reduce framework to use the append operation in the "reduce" phase of the application. To validate this work, we perform experiments on a large number of nodes of the Grid'5000 testbed. We demonstrate that massively concurrent append and read operations have a low impact on each other. Besides, measurements with an application available with Hadoop show that the support for concurrent appends to shared file is introduced with no extra cost, whereas the number of files managed by the Map/Reduced framework is substantially reduced. Diana Moise, Gabriel Antoniu, Luc Bougé |
HPDC | 2 |
| 2010 | BlobSeer: Bringing high throughput under heavy concurrency to Hadoop Map-Reduce applicationsabstractHadoop is a software framework supporting the Map-Reduce programming model. It relies on the Hadoop Distributed File System (HDFS) as its primary storage system. The efficiency of HDFS is crucial for the performance of Map-Reduce applications. We substitute the original HDFS layer of Hadoop with a new, concurrency-optimized data storage layer based on the BlobSeer data management service. Thereby, the efficiency of Hadoop is significantly improved for data-intensive Map-Reduce applications, which naturally exhibit a high degree of data access concurrency. Moreover, BlobSeer's features (built-in versioning, its support for concurrent append operations) open the possibility for Hadoop to further extend its functionalities. We report on extensive experiments conducted on the Grid'5000 testbed. The results illustrate the benefits of our approach over the original HDFS-based implementation of Hadoop. Bogdan Nicolae, Diana Moise, Gabriel Antoniu, Luc Bougé, Matthieu Dorier |
IPDPS | 3 |
| 2009 | Enabling High Data Throughput in Desktop Grids through Decentralized Data and Metadata Management: The BlobSeer Approach
Bogdan Nicolae, Gabriel Antoniu, Luc Bougé |
Euro-Par | 2 |
| 2008 | Enabling lock-free concurrent fine-grain access to massive distributed data: Application to supernovae detectionabstractWe consider the problem of efficiently managing massive data in a large-scale distributed environment. We consider data strings of size in the order of Terabytes, shared and accessed by concurrent clients. On each individual access, a segment of a string, of the order of Megabytes, is read or modified. Our goal is to provide the clients with efficient fine-grain access the data string as concurrently as possible, without locking the string itself. This issue is crucial in the context of applications in the field of astronomy, databases, data mining and multimedia. We illustrate these requirements with the case of an application for searching supernovae. Our solution relies on distributed, RAM-based data storage, while leveraging a DHT-based, parallel metadata management scheme. The proposed architecture and algorithms have been validated through a software prototype and evaluated in a cluster environment. Bogdan Nicolae, Gabriel Antoniu, Luc Bougé |
CLUSTER | 2 |
| 2008 | Building Hierarchical Grid Storage Using the GfarmGlobal File System and the JuxMemGrid Data-Sharing Service
Gabriel Antoniu, Loïc Cudennec, Majd Ghareeb, Osamu Tatebe |
Euro-Par | 1 |
| 2008 | A practical example of convergence of P2P and grid computing: An evaluation of JXTA's communication performance on grid networking infrastructuresabstractAs the size of today's grid computing platforms increases, the need for self-organization and dynamic reconfiguration becomes more and more important. In this context, the convergence of grid computing and peer-to-peer (P2P) computing seems natural. However, grid infrastructures are generally available as a hierarchical federation of SAN-based clusters interconnected by high-bandwidth WANs. In contrast, P2P systems usually run on the Internet, on top of random, generally flat network topologies. This difference may lead to the legitimate question of how adequate are the P2P communication mechanisms on hierarchical grid infrastructures. Answering this question is important, since it is essential to efficiently exploit the particular features of grid networking topologies in order to meet the constraints of scientific applications. This paper evaluates the communication performance of the JXTA P2P platform over high-performance SANs and WANs, for both J2SE and C bindings. We discuss these results, then we propose and evaluate several techniques able to improve the JXTA's performance on such grid networking infrastructures. Gabriel Antoniu, Mathieu Jan, David A. Noblet |
IPDPS | 1 |
| 2007 | Towards a Transparent Data Access Model for the GridRPCParadigm
Gabriel Antoniu, Eddy Caron, Frédéric Desprez, Aurélia Fèvre, Mathieu Jan |
HiPC | 1 |
| 2007 | Performance scalability of the JXTA P2P frameworkabstractFeatures of the P2P model, such as scalability and volatility tolerance, have motivated its use in distributed systems. Several generic P2P libraries have been proposed for building distributed applications. However, very few experimental evaluations of these frameworks have been conducted, especially at large scales. Such experimental analyses are important, since they can help system designers to optimize P2P protocols and better understand the benefits of the P2P model. This is particularly important when the P2P model is applied to special use cases, such as grid computing. This paper focuses on the scalability of two main protocols proposed by the JXTA P2P platform. First, we provide a detailed description of the underlying mechanisms used by JXTA to manage its overlay and propagate messages over it: the rendezvous protocol. Second, we describe the discovery protocol used to find resources inside a JXTA network. We then report a detailed, large-scale, multi-site experimental evaluation of these protocols, using the nine clusters of the French Grid'5000 testbed. Gabriel Antoniu, Loïc Cudennec, Mathieu Jan, Mike Duigou |
IPDPS | 1 |
| 2007 | MOve: Design and Evaluation of a Malleable Overlay for Group-Based ApplicationsabstractWhile peer-to-peer overlays allow distributed applications to scale and tolerate failures, most structured and unstructured overlays in literature today are inflexible from the application viewpoint. The application thus has no first-class control on the overlay structure. This paper proposes the concept of an application-malleable overlay, and the design of the first malleable overlay which we call MOve. MOve is targeted at group- based applications, e.g., collaborative applications. In MOve, the communication characteristics of the distributed application using the overlay can influence the overlay's structure itself, with the twin goals of (1) optimizing the application performance by adapting the overlay, while also (2) retaining the large scale and fault tolerance of the overlay approach. Besides neighbor list membership management, MOve also contains algorithms for resource discovery, update propagation, and churn-resistance. The emergent behavior of the implicit mechanisms used in MOve manifests as follows: when application communication is low, most overlay links keep their default configuration; however, as application communication characteristics become more evident, the overlay gracefully adapts itself to the application. We validate MOve using simulations with group sizes that are fixed, uniform, exponential and PlanetLab-based (slices), as well as churn traces and two sample management-based applications. Ramsés Morales, Sébastien Monnet, Indranil Gupta, Gabriel Antoniu |
IEEE Trans. Netw. Serv. Manag. | 4 |
| 2006 | Extending the Entry Consistency Model to Enable Efficient Visualization for Code-Coupling Grid ApplicationsabstractThis paper addresses the problem of efficient visualization of shared data within code coupling grid applications. These applications are structured as a set of distributed, autonomous, weakly-coupled codes. We focus on the case where the codes are able to interact using the abstraction of a shared data space. We propose an efficient visualization scheme by adapting the mechanisms used to maintain the data consistency. We introduce a new operation called relaxed read, as an extension to the entry consistency model. This operation can efficiently take place without locking, in parallel with write operations. We discuss the benefits and the constraints of the proposed approach. Gabriel Antoniu, Loïc Cudennec, Sébastien Monnet |
CCGRID | 1 |
| 2006 | Enabling Transparent Data Sharing in Component ModelsabstractThe fast growth of high-bandwidth unde-area networks has encouraged the development of computational grids. To deal with the increasing complexity of grid applications, the software component technology seems very appealing since it emphasizes software composition and re-use. However, current software component models only support explicit data transfers between components. The distributed shared memory paradigm has demonstrated its utility by enabling a transparent access to data via a globally shared data space. This paper proposes to extend software component models with shared memory capabilities, enabling tmnsparent access to shared data across components and leading to further decreased software complexity. Gabriel Antoniu, Hinde-Lilia Bouziane, Landry Breuil, Mathieu Jan, Christian Pérez |
CCGRID | 1 |
| 2006 | MOve: Design of An Application-Malleable OverlayabstractPeer-to-peer overlays allow distributed applications to work in a wide-area, scalable, and fault-tolerant manner. However, most structured and unstructured overlays present in literature today are inflexible from the application viewpoint. In other words, the application has no control over the structure of the overlay itself. This paper proposes the concept of an application-malleable overlay, and the design of the first malleable overlay which we call MOve. In MOve, the communication characteristics of the distributed application using the overlay can influence the overlay's structure itself, with the twin goals of (1) optimizing the application performance by adapting the overlay, while also (2) retaining the large scale and fault tolerance of the overlay approach. The influence could either be explicitly specified by the application or implicitly gleaned by our algorithms. Besides neighbor list membership management, MOve also contains algorithms for resource discovery, update propagation, and churn-resistance. The emergent behavior of the implicit mechanisms used in MOve manifests in the following way: when application communication is low, most overlay links keep their default configuration; however, as application communication characteristics become more evident, the overlay gracefully adapts itself to the application Sébastien Monnet, Ramsés Morales, Gabriel Antoniu, Indranil Gupta |
SRDS | 3 |
| 2006 | How to bring together fault tolerance and data consistency to enable Grid data sharingabstractAbstract This paper addresses the challenge of transparent data sharing within computing Grids built as cluster federations. On such platforms, the availability of storage resources may change in a dynamic way, often due to hardware failures. We focus on the problem of handling the consistency of replicated data in the presence of failures. We propose a software architecture which decouples consistency management from fault tolerance management. We illustrate this architecture with a case study showing how to design a consistency protocol using fault‐tolerant building blocks. As a proof of concept, we describe a prototype implementation of this protocol within JUXMEM, a software experimental platform for Grid data sharing, and we report on a preliminary experimental evaluation of the proposed approach. Copyright © 2006 John Wiley & Sons, Ltd. Gabriel Antoniu, Jean-François Deverge, Sébastien Monnet |
Concurr. Comput. Pract. Exp. | 1 |
| 2005 | Performance evaluation of JXTA communication layersabstractThe arrival of the P2P model has opened many new avenues for research within the field of distributed computing. This is mainly due to important practical features (such as support for volatility, high scalability). Several generic P2P libraries have been proposed for building higher-level services. In order to judge the appropriateness of using a generic P2P library for a given application type, an experimental performance evaluation of the provided functionalities is unavoidable. Very few analyses of this kind have been reported, as most evaluations are limited to complexity analyses and to simulations. Such experimental analyses are important, especially when using P2P software in a grid computing context, where applications may have precise efficiency requirements. In this paper, we focus on JXTA, which provides generic building blocks and protocols intended to serve as a basis for specialized P2P services and applications. We perform a performance evaluation of the three communication layers (endpoint, pipe and socket) over a fast Ethernet local-area network, for recent versions of the J2SE and C bindings of JXTA. We provide a detailed analysis explaining the behavior of these three layers and we give hints showing how to efficiently use them. Gabriel Antoniu, Philip J. Hatcher, Mathieu Jan, David A. Noblet |
CCGRID | 1 |
| 2005 | Enabling the P2P JXTA Platform for High-Performance Networking Grid Infrastructures
Gabriel Antoniu, Mathieu Jan, David A. Noblet |
HPCC | 1 |
| 2004 | Large-Scale Deployment in P2P Experiments Using the JXTA Distributed Framework
Gabriel Antoniu, Luc Bougé, Mathieu Jan, Sébastien Monnet |
Euro-Par | 1 |
| 2003 | Making a DSM Consistency Protocol Hierarchy-Aware: an Efficient Synchronization SchemeabstractWe consider the design of DSM consistency protocols for hierarchical architectures. Such architectures typically consist of a constellation of loosely-interconnected clusters, each cluster consisting of a set of tightly-interconnected nodes running multithreaded programs. We claim that high performance can only be reached by taking into account this interconnection hierarchy at the very core of the protocol design. Previous work has focused on improving locality in data management by caching remote data within clusters. In contrast, our idea is to improve locality in the synchronization management. We demonstrate the feasibility through an experimental implementation of this idea in a home-based protocol for Release Consistency, and we provide a preliminary evaluation of the expectable performance gain. Gabriel Antoniu, Luc Bougé, Sébastien Lacour |
CCGRID | 1 |
| 2001 | DSM-PM2: A Portable Implementation Platform for Multithreaded DSM Consistency Protocols
Gabriel Antoniu, Luc Bougé |
HIPS | 1 |
| 2001 | DSM-PM2: A portable implementation platform for multithreaded DSM consistency protocols
Gabriel Antoniu, Luc Bougé |
IPDPS | 1 |
| 2001 | Remote Object Detection in Cluster-Based JavaabstractInternational audience Gabriel Antoniu, Philip J. Hatcher |
IPDPS | 1 |
| 2001 | The Hyperion system: Compiling multithreaded Java bytecode for distributed execution
Gabriel Antoniu, Luc Bougé, Philip J. Hatcher, Mark MacBeth, Keith McGuigan, Raymond Namyst |
Parallel Comput. | 1 |
| 2000 | Compiling Multithreaded Java Bytecode for Distributed Execution (Distinguished Paper)
Gabriel Antoniu, Luc Bougé, Philip J. Hatcher, Mark MacBeth, Keith McGuigan, Raymond Namyst |
Euro-Par | 1 |
| 1999 | Using Preemptive Thread Migration to Load-Balance Data-Parallel Applications
Gabriel Antoniu, Christian Pérez |
Euro-Par | 1 |