VLDB 2026 Research / reviewers in the wild / expert
Alexandru Costan
dblp:81/2916
· DBLP profile ↗
47ranked-venue papers
4as first author
13since 2021 · last 2026
0000-0003-3111-6308ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 27 · 2 first-author · 12 since 2021Artificial intelligence and machine learning · 7Databases, data management, data science and information retrieval · 6Applied, interdisciplinary, general and emerging computing · 6 · 1 since 2021Software engineering, systems software and programming languages · 2 · 1 first-author · 1 since 2021Security and privacy · 1Graphics, computer vision, multimedia, augmented reality and games · 1Theory of computation · 1 · 1 first-author
| 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 | 5 |
| 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 | 3 |
| 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. | 3 |
| 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 | 4 |
| 2024 | Workflow Provenance in the Computing Continuum for Responsible, Trustworthy, and Energy-Efficient AIabstractAs Artificial Intelligence (AI) becomes more pervasive in our society, it is crucial to develop, deploy, and assess Responsible and Trustworthy AI (RTAI) models, i.e., those that consider not only accuracy but also other aspects, such as explainability, fairness, and energy efficiency. Workflow provenance data have historically enabled critical capabilities towards RTAI. Provenance data derivation paths contribute to responsible workflows through transparency in tracking artifacts and resource consumption. Provenance data are well-known for their trustworthiness helping explainability, reproducibility, and accountability. However, there are complex challenges to achieve RTAI, which are further complicated by the heterogeneous infrastructure in the computing continuum (Edge-Cloud-HPC) used to develop and deploy models. As a result, a significant research and development gap remains between workflow provenance data management and RTAI. In this paper, we present a vision of the pivotal role of workflow provenance in supporting RTAI and discuss related challenges. We present a schematic view between RTAI and provenance, and highlight open research directions. Renan Souza 0001, Silvina Caíno-Lores, Mark Coletti, Tyler J. Skluzacek, Alexandru Costan, Frédéric Suter, Marta Mattoso, Rafael Ferreira da Silva |
e-Science | 5 |
| 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 | 3 |
| 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. | 2 |
| 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 | 4 |
| 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 | 3 |
| 2022 | FlexScience'22: 12th Workshop on AI and Scientific Computing at Scale using Flexible Computing InfrastructuresabstractScientific computing applications generate enormous datasets that are continuously increasing exponentially in both complexity and volume, making their analysis, archival, and sharing one of the grand challenges of modern big data analytics. Supported by the rise of artificial intelligence and deep learning, such enormous datasets are becoming valuable resources even beyond their original scope, opening new opportunities to learn patterns and extract new knowledge at large scale, potentially without human intervention. However, this leads to an increasing complexity of the workflows that combine traditional HPC simulations with big data analytics and AI applications. An initial wave that opened this direction was the shift from compute-intensive to data-intensive, which saw several ideas from big data analytics (in-situ processing, shipping computations close to data, complex and dynamic workflows) fused with the tightly coupled patterns addressed by the AI and the high performance computing ecosystems. In a quest to keep up with the complexity of the workflows, the design and operation of the infrastructures capable of running them efficiently at scale has evolved accordingly. Extreme heterogeneity at all levels (combinations of CPUs and accelerators, various types of memories and local storage and network links, parallel file systems and object stores, etc.) is now the norm. ideas pioneered by cloud and edge computing (aspects related to elasticity, multi-tenancy, geo-distributed processing, stream computing) are also beginning to be adopted in the HPC ecosystem (containerized workflows, on-demand jobs to complement batch jobs, streaming of experimental data from instruments directly to supercomputers, etc.). Thus, modern scientific applications need to be integrated into an entire Compute Continuum from the edge all the way to supercomputers and large data-centers using flexible infrastructures and middlewares. The 12th workshop on AI and Scientific Computing at Scale using Flexible Computing Infrastructures (FlexScience) will provide the scientific community a dedicated forum for discussing new research, development, and deployment efforts in running scientific computing workloads in such flexible ecosystems, across the Computing Continuum, focusing on emerging technologies and new convergence challenges that are not sufficiently addressed by the current generation of supercomputers and dedicated data centers. The workshop aims to address questions such as: what architectural changes to existing frameworks (hardware, operating systems, networking and/or programming models) are needed to support flexible computing? Dynamic information derived from remote instruments, coupled simulations, and sensor ensembles that stream data for real-time analysis and machine learning are important emerging trends. How can we leverage and adapt to these patterns? What scientific workloads are suitable candidates to take advantage of heterogeneity, elasticity and/or on-demand resources? What factors are limiting the adoption of a flexible design? Alexandru Costan, Bogdan Nicolae, Kento Sato |
HPDC | 1 |
| 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. | 2 |
| 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 | 2 |
| 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 | 2 |
| 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 | 7 |
| 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 | 4 |
| 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. | 5 |
| 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 | 2 |
| 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. | 4 |
| 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 | 3 |
| 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 | 2 |
| 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 | 4 |
| 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. | 3 |
| 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 | 2 |
| 2017 | A performance evaluation of Apache Kafka in support of big data streaming applicationsabstractStream computing is becoming a more and more popular paradigm as it enables the real-time promise of data analytics. Apache Kafka is currently the most popular framework used to ingest the data streams into the processing platforms. However, how to tune Kafka and how much resources to allocate for it remains a challenge for most users, who now rely mainly on empirical approaches to determine the best parameter settings for their deployments. In this poster, we make a through evaluation of several configurations and performance metrics of Kafka in order to allow users avoid bottlenecks, reach its full potential and avoid bottlenecks and eventually leverage some good practice for efficient stream processing. Paul Le Noac'h, Alexandru Costan, Luc Bougé |
IEEE BigData | 2 |
| 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 | 4 |
| 2017 | AutoCompBD: Autonomic Computing and Big Data platforms
Florin Pop, Ciprian Dobre, Alexandru Costan |
Soft Comput. | 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 | 3 |
| 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 | 2 |
| 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 | 2 |
| 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. | 1 |
| 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. | 2 |
| 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. | 2 |
| 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 | 2 |
| 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 | 2 |
| 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 | 2 |
| 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 | 2 |
| 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 | 4 |
| 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 | 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 | 4 |
| 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 | 3 |
| 2010 | Automatic Generation of Functional Workflows Using a Semantic SpecificationabstractWorkflows are becoming an increasingly more common paradigm to manage and control scientific applications. They are an effective technology to define composition of different pieces of knowledge, both in the application domain and in the scientific context. However, workflows are addressed in a number of different perspectives and little consensus has been reached yet on the specification, languages, modelling and systems to adopt in advanced contexts like Grids or semantic web. In this paper, we describe a tool for supporting users in specifying their functional workflows using a top level description, close to the application domain, and further authoring workflows running on Grids. Our tool is the semantic component of a workflow management platform targeted at scientific applications and is able to automatically generate BPEL processes based on the information extracted from a given ontology. The component has been successfully tested on image processing workflows, using Active BPEL engine as a deployment environment and OWL for ontology design. The functionality provided by the semantic component has been made available through a web interface. Costin Ionita, Alexandru Costan, Valentin Cristea |
CISIS | 2 |
| 2010 | Fault Tolerance and Recovery in Grid Workflow Management SystemsabstractComplex scientific workflows are now commonly executed on global grids. With the increasing scale complexity, heterogeneity and dynamism of grid environments the challenges of managing and scheduling these workflows are augmented by dependability issues due to the inherent unreliable nature of large-scale grid infrastructure. In addition to the traditional fault tolerance techniques, specific checkpoint-recovery schemes are needed in current grid workflow management systems to address these reliability challenges. Our research aims to design and develop mechanisms for building an autonomic workflow management system that will exhibit the ability to detect, diagnose, notify, react and recover automatically from failures of workflow execution. In this paper we present the development of a Fault Tolerance and Recovery component that extends the ActiveBPEL workflow engine. The detection mechanism relies on inspecting the messages exchanged between the workflow and the orchestrated Web Services in search of faults. The recovery of a process from a faulted state has been achieved by modifying the default behavior of ActiveBPEL and it basically represents a non-intrusive checkpointing mechanism. We present the results of several scenarios that demonstrate the functionality of the Fault Tolerance and Recovery component, outlining an increase in performance of about 50% in comparison to the traditional method of resubmitting the workflow. Elvin Sindrilaru, Alexandru Costan, Valentin Cristea |
CISIS | 2 |
| 2010 | Prediction of Distributed Systems State Based on Monitoring DataabstractAutonomic behavior has emerged as a solution for the issues related to performance improvement and resource-usage optimization in large scale distributed systems. This solution relies on monitoring services to keep track of the states of the managed systems. However, most of the monitoring services are designed to provide general resource information and do not consider specific information for higher-level services, lacking important control capabilities. In this context, a dynamic adaptation layer is required. Based on the collected monitoring information in conjunction with some planning and prediction algorithms, it should be able to reactive and proactive deal with detected or predicted conditions. This paper presents a prediction architecture developed within the MonALISA monitoring framework, providing methods for estimating future values for different parameters on various periods of time. The predictions are used to enhance the self-adaptive behavior of several data intensive applications. Our research was focused on machine learning algorithms correlated with statistical techniques for data mining purposes in order to perform n-step-ahead time series predictions and to evaluate their performances dynamical. Adriana Draghici, Alexandru Costan, Valentin Cristea |
ISPDC | 2 |
| 2009 | Dynamic Meta-Scheduling Architecture Based on Monitoring in Distributed SystemsabstractThe Scheduling process in Large Scale Distributed System (LSDS) became more important due of increases of users and applications. This paper presents a dynamic meta-scheduling architecture model for LSDS based on monitoring. Dynamic scheduling process tries to perform task allocation on the fly as the application executes. The monitoring is important in this process because can offer a full view of nodes in distributed systems. The proposed architecture is an agent framework and contains a Grid Monitoring Service, an Execution Services and a Discovery Services. The performance of used monitoring system-MonALISA is very important for dynamic scheduling because ensure the real-time process. The experimental results validate our architecture and scheduling model. Florin Pop, Ciprian Dobre, Corina Stratan, Alexandru Costan, Valentin Cristea |
CISIS | 4 |
| 2009 | Monitoring of Complex Applications Execution in Distributed Dependable SystemsabstractThe execution of applications in dependable system requires a high level of instrumentation for automatic control. We present in this paper a monitoring solution for complex application execution. The monitoring solution is dynamic, offering real-time information about systems and applications. The complex applications are described using workflows. We show that the management process for application execution is improved using monitoring information. The environment is represented by distributed dependable systems that offer a flexible support for complex application execution. Our experimental results highlight the performance of the proposed monitoring tool, the MonALISA framework. Florin Pop, Alexandru Costan, Ciprian Dobre, Corina Stratan, Valentin Cristea |
ISPDC | 2 |
| 2008 | A Monitoring Architecture for High-Speed Networks in Large Scale Distributed CollaborationsabstractIn this paper we present the architecture of a distributed framework that allows real-time accurate monitoring of large scale high-speed networks. An important component of a large-scale distributed collaboration is the complex network infrastructure on which it relies. For monitoring and controlling the networking resources an adequate instrument should offer the possibility to collect and store the relevant monitoring information, presenting significant perspectives and synthetic views of how the large distributed system performs. We therefore developed within the MonALISA monitoring framework a system able to collect, store, process and interpret the large volume of status information related to the US LHCNet research network. The system uses flexible mechanisms for data representation, providing access optimization and decision support, being able to present real-time and long-time history information through global or specific views and to take further automated control actions based on them. Alexandru Costan, Ciprian Dobre, Valentin Cristea, Ramiro Voicu |
ISPDC | 1 |
| 2005 | A Policy Iteration Algorithm for Computing Fixed Points in Static Analysis of Programs
Alexandru Costan, Stéphane Gaubert, Eric Goubault, Matthieu Martel, Sylvie Putot |
CAV | 1 |