Alexandru Iosup

dblp:85/2489 · DBLP profile ↗
← Back
117ranked-venue papers
21as first author
28since 2021 · last 2026
0000-0001-8030-9398ORCID · verified

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

Systems, architecture and hardware · 71 · 14 first-author · 17 since 2021Software engineering, systems software and programming languages · 23 · 2 first-author · 10 since 2021Computer networks · 10Graphics, computer vision, multimedia, augmented reality and games · 5Databases, data management, data science and information retrieval · 4 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 4 · 1 first-author · 1 since 2021Artificial intelligence and machine learning · 3Human-computer interaction and ubiquitous computing · 2 · 1 first-author
YearPublicationVenuePosition
2026 OpenDC-STEAM: Realistic Modeling and Systematic Exploration of Composable Techniques for Sustainable Datacenters
abstract
The need to reduce datacenter carbon-footprint is urgent. While many sustainability techniques have been proposed, they are often evaluated in isolation, using limited setups or analytical models that overlook real-world dynamics and interactions between methods. This makes it challenging for researchers and operators to understand the effectiveness and trade-offs of combining such techniques. We design OpenDC-STEAM, an open-source customizable datacenter simulator, to investigate the individual and combined impact of sustainability techniques on datacenter operational and embodied carbon emissions, and their trade-off with performance. Using STEAM, we systematically explore three representative techniques-horizontal scaling, leveraging batteries, and temporal shifting-with diverse representative workloads, datacenter configurations, and carbon-intensity traces. Our analysis highlights that datacenter dynamics can influence their effectiveness and that combining strategies can significantly lower emissions, but introduces complex cost-emissions-performance trade-offs that STEAM can help navigate. STEAM supports the integration of new models and techniques, making it a foundation framework for holistic, quantitative, and reproducible research in sustainable computing. Following open-science principles, STEAM is available as FOSS: https://github.com/atlarge-research/OpenDC-STEAM.
Dante Niewenhuis, Sacheendra Talluri, Alexandru Iosup, Tiziano De Matteis
CCGrid3
2026 M3SA: Exploring Datacenter Performance and Climate-Impact with Multi- and Meta-Model Simulation and Analysis
abstract
Datacenters are vital for the digital society but represent a considerable fraction of global energy consumption. To improve their sustainability and performance when demand is foreseen to increase, we envision simulators will become primary decision-making tools. However, unlike other fields focusing on key societal infrastructure such as waterworks and mass transit, datacenter simulators cannot yet combine multiple, independent models into their operation. Addressing this challenge, in this work we propose M3SA, a datacenter simulation and analysis framework that uses discrete-event simulation to predict, per model and then combined into a meta-model, the impact on climate and performance of various realistic datacenter conditions. We design an architecture for simulating multiple concurrent models, a technique to integrate the results of multiple models into a meta-model, and a procedure to evaluate the accuracy of the meta-model. Through experiments with a prototype, we show that (i) M3SA can be used to reproduce peer-reviewed experiments, and enhance their output with more diverse metrics and more detailed analysis; (ii) M3SA can be configured with a variety of realistic parameters, such as diverse workload traces (using Grid Workloads Archive data in our experiments), and energy production data such as carbon intensity over time and location (gCO2/kWh data across all EU regions, from ENTSO-E); (iii) M3SA enables various types of what-if and how-to analysis, such as how to configure CO2-aware migration over yearly energy-production patterns. M3SA has been integrated into the open-source software simulator OpenDC and is available on https://github.com/atlarge-research/opendc-m3sa.
Radu Nicolae, Dante Niewenhuis, Sacheendra Talluri, Alexandru Iosup
CF4
2026 Let's trace it: Fine-grained serverless benchmarking for synchronous and asynchronous applications
abstract
Making serverless computing widely applicable requires detailed understanding of performance. Although benchmarking approaches exist, their insights are coarse-grained and typically insufficient for (root cause) analysis of realistic serverless applications, which often consist of asynchronously coordinated functions and services. Addressing this gap, we design and implement ServiTrace, an approach for fine-grained distributed trace analysis and an application-level benchmarking suite for diverse serverless-application architectures. ServiTrace (i) analyzes distributed serverless traces using a novel algorithm and heuristics for extracting a detailed latency breakdown , (ii) leverages a suite of serverless applications representative of production usage, including synchronous and asynchronous serverless applications with external service integrations, and (iii) automates comprehensive, end-to-end experiments to capture application-level performance. Using our ServiTrace reference implementation, we conduct a large-scale empirical performance study in the market-leading AWS environment, collecting over 7.5 million execution traces. We make four main observations enabled by our latency breakdown analysis of median latency, cold starts, and tail latency for different application types and invocation patterns. For example, the median end-to-end latency of serverless applications is often dominated not by function computation but by external service calls, orchestration, and trigger-based coordination; all of which could be hidden without ServiTrace-like benchmarking. We release empirical data under FAIR principles and ServiTrace as a tested, extensible, open-source tool at https://github.com/ServiTrace/ReplicationPackage .
Joel Scheuner, Simon Eismann, Sacheendra Talluri, Erwin Van Eyk, Cristina L. Abad, Philipp Leitner 0001, Alexandru Iosup
Future Gener. Comput. Syst.7
2026 Cloud Uptime Archive: Open-Access Availability Data of Web, Cloud, and Gaming Services
abstract
Cloud services are critical to society. However, their reliability is poorly understood. Towards solving the problem, we propose a standard repository for cloud uptime data. We populate this repository with the data we collect containing failure reports from users and operators of cloud services, web services, and online games. The multiple vantage points help reduce bias from individual users and operators. We compare our new data to existing failure data from the Failure Trace Archive and the Google cluster trace. We analyze the MTBF and MTTR, time patterns, failure severity, user-reported symptoms, and operator-reported symptoms of failures in the data we collect. We observe that high-level user facing services fail less often than low-level infrastructure services, likely due to them using fault-tolerance techniques. We use simulation-based experiments to demonstrate the impact of different failure traces on the performance of checkpointing and retry mechanisms. We release the data, and the analysis and simulation tools, as open-source artifacts available athttps://github.com/atlarge-research/cloud-uptime-archive.
Sacheendra Talluri, Dante Niewenhuis, Xiaoyu Chu, Jakob Kyselica, Mehmet Çetin, Alexander Balgavy, Alexandru Iosup
IEEE Trans. Parallel Distributed Syst.7
2025 Performance Characterization of Data Store Event Trigger Mechanisms for Serverless Computing
abstract
Serverless applications are composed of functions triggered by events. Data stores are a common source of event triggers in the cloud, even beyond serverless, such as in Kubernetes. We find trigger latency, the time from event generation to function invocation, to take up to 62% of execution time for common serverless applications. Even though event triggers play a crucial role in serverless performance, the mechanisms driving these triggers are ill-understood. In this paper, we analyze data store trigger mechanisms, define the features that make up these mechanisms, and characterize their performance with TriggerPerf, a benchmarking tool for data store triggers. We implement TriggerPerf on three AWS data stores with built-in trigger support: S3, DynamoDB, and AuroraDB. With TriggerPerf, we demonstrate significant latency, scalability, and elasticity bottlenecks across these data stores. We observe that the trigger latency of AWS data stores is up to$100 \times$higher compared to a reference etcd data store. Moreover, the median tail latency of S3 and AuroraDB is 10x higher when under high load, unlike DynamoDB. The observed variability in performance patterns significantly impacts the reliability of serverless and distributed systems that depend on them, highlighting the critical need for further research into the underlying mechanisms. The tool is open-sourced and is available at https://github.com/atlarge-research/trigger-perf.
Ritul Satish, Sacheendra Talluri, Sudarsan Sivakumar, Matthijs Jansen, Alexandru Iosup
CCGrid5
2025 An Empirical Characterization of Outages and Incidents in Public Services for Large Language Models
abstract
People and businesses increasingly rely on public LLM services, such as ChatGPT, DALL·E, and Claude. Understanding their outages, and particularly measuring their failure-recovery processes, is becoming a stringent problem. However, only limited studies exist in this emerging area. Addressing this problem, in this work we conduct an empirical characterization of outages and failure-recovery in public LLM services. We collect and prepare datasets for 8 commonly used LLM services across 3 major LLM providers, including market-leads OpenAI and Anthropic. We conduct a detailed analysis of failure recovery statistical properties, temporal patterns, co-occurrence, and the impact range of outage-causing incidents. We make over 10 observations, among which: (1) Failures in OpenAI's ChatGPT take longer to resolve but occur less frequently than those in Anthropic's Claude;(2) OpenAI and Anthropic service failures exhibit strong weekly and monthly periodicity; and (3) OpenAI services offer better failure-isolation than Anthropic services. Our research explains LLM failure characteristics and thus enables optimization in building and using LLM systems. FAIR data and code are publicly available on https://zenodo.org/records/14018219 and https://github.com/atlarge-research/llm-service-analysis.
Xiaoyu Chu, Sacheendra Talluri, Qingxian Lu, Alexandru Iosup
ICPE4
2025 Columbo: A Reasoning Framework for Kubernetes' Configuration Space
abstract
Resource managers such as Kubernetes are rapidly evolving to support low-latency and scalable computing paradigms such as serverless and granular computing. As a result, Kubernetes supports dozens of workload deployment models and exposes roughly 1,600 configuration parameters. Previous work has shown that parameter tuning can significantly improve Kubernetes' performance, but identifying which parameters impact performance and should be tuned remains challenging. To help users optimize their Kubernetes deployments, we present Columbo, an offline reasoning framework to detect and resolve performance bottlenecks using configuration parameters. We study Kubernetes and define its workload deployment pipeline of 6 stages and 26 steps. To detect bottlenecks, Columbo uses an analytical model to predict the best-case deployment time of a workload per pipeline stage and compares it to empirical data from a novel benchmark suite. Columbo then uses a rule-based methodology to recommend parameter updates based on the detected bottleneck, deployed workload, and mapping of configurations to pipeline stages. We demonstrate that Columbo reduces workload deployment time across its benchmark suite by 28% on average and 79% at most. We report a total execution time decrease of 17% for data processing with Spark and up to 20% for serverless workflows with OpenWhisk. Columbo is open-source and available at https://github.com/atlarge-research/continuum/tree/columbo.
Matthijs Jansen, Sacheendra Talluri, Krijn Doekemeijer, Nick Tehrany, Alexandru Iosup, Animesh Trivedi
ICPE5
2025 RADiCe: A Risk Analysis Framework for Data Centers
abstract
Datacenter service providers face engineering and operational challenges involving numerous risk aspects. Bad decisions can result in financial penalties, competitive disadvantage, and unsustainable environmental impact. Risk management is an integral aspect of the design and operation of modern datacenters, but frameworks that allow users to consider various risk trade-offs conveniently are missing. We propose RADICE, an open-source framework that enables data-driven analysis of IT-related operational risks in sustainable datacenters. RADICE uses monitoring and environmental data and, via discrete event simulation, assists datacenter experts through systematic evaluation of risk scenarios, visualization, and optimization of risks. Our analyses highlight the increasing risk datacenter operators face due to price surges in electricity and sustainability and demonstrate how RADICE can evaluate and control such risks by optimizing the topology and operational settings of the datacenter. Eventually, RADICE can evaluate risk scenarios by a factor 70x–330x faster than others, opening possibilities for interactive risk exploration.
Fabian Mastenbroek, Tiziano De Matteis, Vincent van Beek, Alexandru Iosup
Future Gener. Comput. Syst.4
2025 Serverless Computing for Next-generation Application Development
abstract
Serverless computing is a cloud computing model that abstracts server management, allowing developers to focus solely on writing code without concerns about the underlying infrastructure. This paradigm shift is transforming application development by reducing time to market, lowering costs, and enhancing scalability. In serverless computing, functions are event-driven and automatically scale in response to events such as data changes or user requests. Despite its advantages, serverless computing presents several research challenges, including managing state for ephemeral functions, mitigating cold start delays, optimizing function composition, debugging, efficient auto-scaling, resource management, and ensuring security and compliance. This special issue focused on addressing these challenges by promoting research on innovative solutions and exploring the potential of serverless computing in new application domains.
Adel Nadjaran Toosi, Bahman Javadi, Alexandru Iosup, Evgenia Smirni, Schahram Dustdar
Future Gener. Comput. Syst.3
2024 Generic and ML Workloads in an HPC Datacenter: Node Energy, Job Failures, and Node-Job Analysis
abstract
HPC datacenters offer a backbone to the modern digital society. Increasingly, they run Machine Learning (ML) jobs next to generic, compute-intensive workloads, supporting science, business, and other decision-making processes. However, understanding how ML jobs impact the operation of HPC datacenters, relative to generic jobs, remains desirable but understudied. In this work, we leverage long-term operational data, collected from a national-scale production HPC datacenter, and statistically compare how ML and generic jobs can impact the performance, failures, resource utilization, and energy consumption of HPC datacenters. Our study provides key insights, e.g., ML-related power usage causes GPU nodes to run into temperature limitations, median/mean runtime and failure rates are higher for ML jobs than for generic jobs, both ML and generic jobs exhibit highly variable arrival processes and resource demands, significant amounts of energy are spent on unsuccessfully terminating jobs, and concurrent jobs tend to terminate in the same state. We open-source our cleaned-up data traces on Zenodo (https://doi. org/10.5281/zenodo.13685426), and provide our analysis toolkit as software hosted on GitHub (https://github.com/atlarge-research/2024-icpads-hpc-workload-characterization). This study offers multiple benefits for data center administrators, who can improve operational efficiency, and for researchers, who can further improve system designs, scheduling techniques, etc.
Xiaoyu Chu, Daniel Hofstätter, Shashikant Ilager, Sacheendra Talluri, Duncan Kampert, Damian Podareanu, Dmitry Duplyakin, Ivona Brandic, Alexandru Iosup
ICPADS9
2024 The Cost of Simplicity: Understanding Datacenter Scheduler Programming Abstractions
abstract
Schedulers are a crucial component in datacenter resource management. Each scheduler offers different capabilities, and users use them through their APIs. However, there is no clear understanding of what programming abstractions they offer, nor why they offer some and not others. Consequently, it is difficult to understand their differences and the performance costs imposed by their APIs. In this work, we study the programming abstractions offered by industrial schedulers, their shortcomings, and their related performance costs. We propose a general reference architecture for scheduler programming abstractions. Specifically, we analyze the programming abstractions of five popular industrial schedulers, understand the differences in their APIs, and identify the missing abstractions. Finally, we carry out exemplary experiments using trace-driven simulation demonstrating that an API extension, such as container migration, can improve total execution time per task by 81%, highlighting how schedulers sacrifice performance by implementing simpler programming abstractions. All the relevant software and data artifacts are publicly available at https://github.com/atlarge-research/quantifying-api-design.
Aratz Manterola Lasa, Sacheendra Talluri, Tiziano De Matteis, Alexandru Iosup
ICPE4
2024 ExDe: Design space exploration of scheduler architectures and mechanisms for serverless data-processing
abstract
Serverless computing is increasingly used for data-processing applications in both science and business domains. At the core of serverless data-processing systems is the scheduler, which ensures dynamic decisions about task and data placement. Due to the variety of user, cluster, and workload properties, the design space for high-performance and cost-effective scheduling architectures and mechanisms is vast. The large design space is difficult to explore and characterize. To help the system designer disentangle this complexity, we present ExDe, a framework to systematically explore the design space of scheduling architectures and mechanisms. The framework includes a conceptual model and a simulator to assist in design space exploration. We use the framework, and real-world workloads, to characterize the performance of three scheduling architectures and two mechanisms. Our framework is open-source software available on Zenodo.
Sacheendra Talluri, Nikolas Herbst, Cristina L. Abad, Tiziano De Matteis, Alexandru Iosup
Future Gener. Comput. Syst.5
2023 The SPEC-RG Reference Architecture for The Compute Continuum
abstract
As the next generation of diverse workloads like autonomous driving and augmented/virtual reality evolves, computation is shifting from cloud-based services to the edge, leading to the emergence of a cloud-edge compute continuum. This continuum promises a wide spectrum of deployment opportunities for workloads that can leverage the strengths of cloud (scalable infrastructure, high reliability) and edge (energy efficient, low latencies). Despite its promises, the continuum has only been studied in silos of various computing models, thus lacking strong end-to-end theoretical and engineering foundations for computing and resource management across the continuum. Consequently, devel-opers resort to ad hoc approaches to reason about performance and resource utilization of workloads in the continuum. In this work, we conduct a first-of-its-kind systematic study of various computing models, identify salient properties, and make a case to unify them under a compute continuum reference architecture. This architecture provides an end-to-end analysis framework for developers to reason about resource management, workload distribution, and performance analysis. We demonstrate the utility of the reference architecture by analyzing two popular continuum workloads, deep learning and industrial IoT. We have developed an accompanying deployment and benchmarking framework and first-order analytical model for quantitative reasoning of continuum workloads. The framework is open-sourced and available at https://github.com/atlarge-research/continuum.
Matthijs Jansen, Auday Aldulaimy, Alessandro Vittorio Papadopoulos, Animesh Trivedi, Alexandru Iosup
CCGrid5
2023 A Trace-driven Performance Evaluation of Hash-based Task Placement Algorithms for Cache-enabled Serverless Computing
abstract
Data-driven interactive computation is widely used for business analytics, search-based decision-making, and log mining. These applications' short duration and bursty nature makes them a natural fit for serverless computing. Data processing serverless applications are composed of many small tasks. Application tasks that use remote storage encounter bottlenecks in the form of high latency, performance variability, and throttling. Caching has been used to mitigate this bottleneck for intermediate data. However, the use of caching for input data, albeit widely used in industry, has yet to be studied. We present the first performance study of scaling, a key feature of serverless computing, on serverless clusters with input data caches. We compare 8 task placement algorithms and quantify their impact on task slowdown and resource usage before and after scaling. We quantify the consequences of using work stealing. We quantify the performance impact of scaling in the buffer period immediately after scaling. We find up to a 420% increase in task slowdown after scaling without work stealing and a 22% slowdown with work stealing. We also find that cache misses after scaling can lead to an additional 21% resource usage.
Sacheendra Talluri, Nikolas Herbst, Cristina L. Abad, Animesh Trivedi, Alexandru Iosup
CF5
2023 Massivizing Computer Systems: VU on the Science, Design, and Engineering of Distributed Systems and Ecosystems
Alexandru Iosup
CLOSER1
2023 The Graph-Massivizer Approach Toward a European Sustainable Data Center Digital Twin
abstract
Modeling and understanding an expensive next-generation data center operating at a sustainable exascale performance remains a challenge yet to solve. The paper presents the approach taken by the Graph-Massivizer project, funded by the European Union, towards a sustainable data center, targeting a massive graph representation and analysis of its digital twin. We introduce five interoperable open-source tools that support this undertaking, creating an automated, sustainable loop of graph creation, analytics, optimization, sustainable resource management, and operation, emphasizing state-of-the-art progress. We plan to employ the tools for designing a massive data center graph, representing a digital twin describing spatial, semantic, and temporal relationships between the monitoring metrics, hardware nodes, cooling equipment, and jobs. The project aims to strengthen Bologna Technopole as a leading European supercomputing and big data hub offering sustainable green computing for improved societally relevant science throughput.
Martin Molan, Junaid Ahmed Khan, Andrea Bartolini, Roberta Turra, Giorgio Pedrazzi, Michael Cochez, Alexandru Iosup, Dumitru Roman, Joze M. Rozanec, Ana Lucia Varbanescu, Radu Prodan
COMPSAC7
2023 Servo: Increasing the Scalability of Modifiable Virtual Environments Using Serverless Computing
abstract
Online games with modifiable virtual environments (MVEs) have become highly popular over the past decade. Among them, Minecraft-supporting hundreds of millions of users―is the best-selling game of all time, and is increasingly offered as a service. Although Minecraft is architected as a distributed system, in production it achieves this scale by partitioning small groups of players over isolated game instances. From the approaches that can help other kinds of virtual worlds scale, none is designed to scale MVEs, which pose a unique challenge―a mix between the count and complexity of active in-game constructs, player-created in-game programs, and strict quality of service. Serverless computing emerged recently and focuses, among others, on service scalability. Thus, addressing this challenge, in this work we explore using serverless computing to improve MVE scalability. To this end, we design, prototype, and evaluate experimentally Servo, a serverless backend architecture for MVEs. We implement Servo as a prototype and evaluate it using real-world experiments on two commercial serverless platforms, of Amazon Web Services (AWS) and Microsoft Azure. Results offer strong support that our serverless MVE can significantly increase the number of supported players per instance without performance degradation, in our key experiment by 40 to 140 players per instance, which is a significant improvement over state-of-the-art commercial and open-source alternatives. We release Servo as open-source, on Github: https://github.com/atlarge-research/opencraft.
Jesse Donkervliet, Javier Ron, Tiberiu Iancu, Cristina L. Abad, Alexandru Iosup
ICDCS6
2023 Meterstick: Benchmarking Performance Variability in Cloud and Self-hosted Minecraft-like Games
abstract
Due to increasing popularity and strict performance requirements, online games have become a workload of interest for the performance engineering community. One of the most popular types of online games is the Minecraft-like Game (MLG), in which players can terraform the environment. The most popular MLG, Minecraft, provides not only entertainment, but also educational support and social interaction, to over 130 million people world-wide. MLGs currently support their many players by replicating isolated instances that support each only up to a few hundred players under favorable conditions. In practice, as we show here, the real upper limit of supported players can be much lower. In this work, we posit that performance variability is a key cause for the lack of scalability in MLGs, investigate experimentally causes of performance variability, and derive actionable insights. We propose a novel operational model for MLGs and use it to design the first benchmark that focuses on MLG performance variability, defining specialized workloads, metrics, and processes. We conduct real-world benchmarking of MLGs, both cloud-based and self-hosted, and find environment-based workloads and cloud deployment to be significant sources of performance variability: peak-latency degrades sharply to 20.7 times the arithmetic mean, and exceeds by a factor of 7.4 the performance requirements. We derive actionable insights for game-developers, game-operators, and other stakeholders to tame performance variability.
Jerrit Eickhoff, Jesse Donkervliet, Alexandru Iosup
ICPE3
2023 Less is not more: We need rich datasets to explore
abstract
Traditional datacenter analysis is based on high-level, coarse-grained metrics. This obscures our vision of datacenter behavior, as we do not observe the full picture nor subtleties that might make up these high-level, coarse metrics. There is room for operational improvement based on fine-grained temporal and spatial, low-level metric data. We leverage in this work one of the (rare) public datasets providing fine-grained information on datacenter operations, with over 60 billion measurements captured in 15-second intervals. We show evidence that fine-grained information reveals new operational aspects, that the different metrics cannot be derived from one another (and thus need to be captured), and that many low-level metrics, gathered frequently are key to understanding datacenter operations. We propose a holistic analysis for datacenter operations, providing statistical characterization of node and workload aspects. Our analysis reveals both generic and machine learning-specific aspects, summarized in over 30 observations, providing deep insight into this dataset and the originating cluster. We give actionable insights, surprising findings, and exemplify how our observations support performance-engineering tasks such as workload prediction and long-term datacenter design.
Laurens Versluis, Mehmet Çetin, Caspar Greeven, Kristian Laursen, Damian Podareanu, Valeriu Codreanu, Alexandru Uta, Alexandru Iosup
Future Gener. Comput. Syst.8
2022 Meterstick: Benchmarking Performance Variability in Cloud and Self-hosted Minecraft-like Games
abstract
One of the most popular types of online games is the Minecraft-like Game (MLG), in which players can terraform the environment. MLGs currently support their many players by replicating isolated instances with limited scalability. We posit that performance variability is a key cause for the lack of scalability in MLGs and design the first benchmark that focuses on MLG performance variability, identifying specialized workloads, metrics, and processes. We conduct real-world benchmarking of MLGs, both cloud-based and self-hosted. We find environment-based workloads and cloud deployment are significant sources of performance variability: peak-latency degrades sharply to 20.7 times the arithmetic mean, and exceeds by a factor of 7.4 the performance requirements.
Jerrit Eickhoff, Jesse Donkervliet, Alexandru Iosup
ISPASS3
2022 Capelin: Data-Driven Compute Capacity Procurement for Cloud Datacenters Using Portfolios of Scenarios
abstract
Cloud datacenters provide a backbone to our digital society. Inaccurate capacity procurement for cloud datacenters can lead to significant performance degradation, denser targets for failure, and unsustainable energy consumption. Although this activity is core to improving cloud infrastructure, relatively few comprehensive approaches and support tools exist for mid-tier operators, leaving many planners with merely rule-of-thumb judgement. We derive requirements from a unique survey of experts in charge of diverse datacenters in several countries. We propose Capelin, a data-driven, scenario-based capacity planning system for mid-tier cloud datacenters. Capelin introduces the notion of portfolios of scenarios, which it leverages in its probing for alternative capacity-plans. At the core of the system, a trace-based, discrete-event simulator enables the exploration of different possible topologies, with support for scaling the volume, variety, and velocity of resources, and for horizontal (scale-out) and vertical (scale-up) scaling. Capelin compares alternative topologies and for each gives detailed quantitative operational information, which could facilitate human decisions of capacity planning. We implement and open-source Capelin, and show through comprehensive trace-based experiments it can aid practitioners. The results give evidence that reasonable choices can be worse by a factor of 1.5-2.0 than the best, in terms of performance degradation or energy consumption.
George Andreadis, Fabian Mastenbroek, Vincent van Beek, Alexandru Iosup
IEEE Trans. Parallel Distributed Syst.4
2022 The State of Serverless Applications: Collection, Characterization, and Community Consensus
abstract
Over the last five years, all major cloud platform providers have increased their serverless offerings. Many early adopters report significant benefits for serverless-based over traditional applications, and many companies are considering moving to serverless themselves. However, currently there exist only few, scattered, and sometimes even conflicting reports on when serverless applications are well suited and what the best practices for their implementation are. We address this problem in the present study about the state of serverless applications. We collect descriptions of 89 serverless applications from open-source projects, academic literature, industrial literature, and domain-specific feedback. We analyze 16 characteristics that describe why and when successful adopters are using serverless applications, and how they are building them. We further compare the results of our characterization study to 10 existing, mostly industrial, studies and datasets; this allows us to identify points of consensus across multiple studies, investigate points of disagreement, and overall confirm the validity of our results. The results of this study can help managers to decide if they should adopt serverless technology, engineers to learn about current practices of building serverless applications, and researchers and platform providers to better understand the current landscape of serverless applications.
Simon Eismann, Joel Scheuner, Erwin Van Eyk, Maximilian Schwinger, Johannes Grohmann, Nikolas Herbst, Cristina L. Abad, Alexandru Iosup
IEEE Trans. Software Eng.8
2021 OpenDC 2.0: Convenient Modeling and Simulation of Emerging Technologies in Cloud Datacenters
abstract
Cloud datacenters are important for the digital society, serving stakeholders across industry, government, and academia. Simulation is a critical part of exploring datacenter technologies, enabling scalable experimentation with millions of jobs and hundreds of thousands of machines, and what-if analysis in a matter of minutes to hours. Although the community has already developed powerful simulators, emerging technologies and applications in modern datacenters require new approaches. Addressing this requirement, in this work we propose OpenDC, a new platform for datacenter simulation. OpenDC includes novel models for emerging cloud-datacenter technologies and applications, such as serverless computing with FaaS deployment and TensorFlow-based machine learning. Our design also focuses on convenience, with a web-based interface for interactive experimentation, support for experiment automation, a library of prefabs for constructing and sharing datacenter designs, and support for diverse input formats and output metrics. We implement, validate, and open-source OpenDC 2.0, a significant redesign and release after a multi-year research and development process. We demonstrate the benefits of OpenDC for the field through a set of representative use-cases: serverless, machine learning, procurement of HPC-as-a-Service infrastructure, educational practices, and reproducibility studies. Overall, OpenDC helps understand how datacenters work, design datacenter infrastructure, and train the next generation of experts.
Fabian Mastenbroek, George Andreadis, Soufiane Jounaid, Wenchen Lai, Jacob Burley, Jaro Bosch, Erwin Van Eyk, Laurens Versluis, Vincent van Beek, Alexandru Iosup
CCGRID10
2021 Dyconits: Scaling Minecraft-like Services through Dynamically Managed Inconsistency
abstract
Gaming is one of the most popular and lucrative entertainment industries. Minecraft alone exceeds 130 million active monthly players and sells millions of licenses annually; it is also provided as a (paid) service. Minecraft, and thousands of others, provide each a Modifiable Virtual Environment (MVE). However, Minecraft-like games only scale using isolated instances that support at most a few hundred players in the same virtual world, thus preventing their large player-base from actually gaming together. When operating as a service, even fewer players can game together. Existing techniques for managing data in distributed systems do not scale for such games: they either do not work for high-density areas (e.g., village centers or other places where the MVE is often modified), or can introduce an unbounded amount of inconsistency that can lower the quality of experience. In this work, we propose Dyconits, a middleware that allows games to scale, by bounding inconsistency in MVEs, optimistically and dynamically. Dyconits allow game developers to partition offline the game-world and its objects into units, each with its own bounds. The Dyconits system controls, dynamically and policy-based, the creation of dyconits and the management of their bounds. Importantly, the Dyconits system is thin, and reuses the existing game codebase and in particular the network stack. To demonstrate and evaluate Dyconits in practice, we modify an existing, open-source, Minecraft-like game, and evaluate its effectiveness through real-world experiments. Our approach supports up to 40% more concurrent players and reduces network bandwidth by up to 85%, with only minor modifications to the game and without increasing game latency.
Jesse Donkervliet, Jim Cuijpers, Alexandru Iosup
ICDCS3
2021 The Fourth Workshop on Hot Topics in Cloud Computing Performance (HotCloudPerf'21): Benchmarking in the Cloud
abstract
The HotCloudPerf workshop is a meeting venue for academics and practitioners, from experts to trainees, in the field of cloud computing performance. The workshop aims to engage this community, and to lead to the development of new methodological aspects for gaining deeper understanding not only of cloud performance, but also of cloud operation and behavior, through diverse quantitative evaluation tools, including benchmarks, metrics, and workload generators. The workshop focuses on novel cloud properties such as elasticity, performance isolation, dependability, and other non-functional system properties, in addition to classical performance-related metrics such as response time, throughput, scalability, and efficiency. The theme for the 2021 edition is "Benchmarking in the Cloud". HotCloudPerf 2021, co-located with the 12th ACM/SPEC International Conference on Performance Engineering (ICPE 2021), is held on April 19-20th, 2021.
Cristina L. Abad, Nikolas Herbst, Alexandru Uta, Alexandru Iosup
ICPE4
2021 Welcome to the 3rd Workshop on Education and Practice of Performance Engineering
abstract
The Workshop on Education and Practice of Performance Engineering, in its 3rd edition, brings together University researchers and Industry Performance Engineers to share education and practice experiences. The recommendations from previous WEPPE workshops pointed to the need to form a team of critical thinkers that also have good communication skills.
Alberto Avritzer, Kishor S. Trivedi, Alexandru Iosup
ICPE3
2021 A survey of domains in workflow scheduling in computing infrastructures: Community and keyword analysis, emerging trends, and taxonomies
abstract
Workflows are prevalent in today’s computing infrastructures as they support many domains. Different Quality of Service (QoS) requirements of both users and providers makes workflow scheduling challenging. Meeting the challenge requires an overview of state-of-art in workflow scheduling. Sifting through literature to find the state-of-art can be daunting, for both newcomers and experienced researchers. Surveys are an excellent way to address questions regarding the different techniques, policies, emerging areas, and opportunities present, yet they rarely take a systematic approach and publish their tools and data on which they are based. Moreover, the communities behind these articles are rarely studied. We attempt to address these shortcomings in this work. We introduce and open-source an instrument used to combine and store article meta-data. Using this meta-data, we characterize and taxonomize the workflow scheduling community and four areas within workflow scheduling: (1) the workflow formalism, (2) workflow allocation, (3) resource provisioning, and (4) applications and services. In each characterization, we obtain important keywords overall and per year, identify keywords growing in importance, get insight into the structure and relations within each community, and perform a systematic literature survey per part to validate and complement our taxonomies
Laurens Versluis, Alexandru Iosup
Future Gener. Comput. Syst.2
2021 Methodological Principles for Reproducible Performance Evaluation in Cloud Computing
abstract
The rapid adoption and the diversification of cloud computing technology exacerbate the importance of a sound experimental methodology for this domain. This work investigates how to measure and report performance in the cloud, and how well the cloud research community is already doing it. We propose a set of eight important methodological principles that combine best-practices from nearby fields with concepts applicable only to clouds, and with new ideas about the time-accuracy trade-off. We show how these principles are applicable using a practical use-case experiment. To this end, we analyze the ability of the newly released SPEC Cloud IaaS benchmark to follow the principles, and showcase real-world experimental studies in common cloud environments that meet the principles. Last, we report on a systematic literature review including top conferences and journals in the field, from 2012 to 2017, analyzing if the practice of reporting cloud performance measurements follows the proposed eight principles. Worryingly, this systematic survey and the subsequent two-round human reviews, reveal that few of the published studies follow the eight experimental principles. We conclude that, although these important principles are simple and basic, the cloud community is yet to adopt them broadly to deliver sound measurement of cloud environments.
Alessandro Vittorio Papadopoulos, Laurens Versluis, André Bauer 0001, Nikolas Herbst, Jóakim von Kistowski, Ahmed Ali-Eldin, Cristina L. Abad, José Nelson Amaral, Petr Tuma 0001, Alexandru Iosup
IEEE Trans. Software Eng.10
2020 Grade10: A Framework for Performance Characterization of Distributed Graph Processing
abstract
Graph processing is one of the most important and ubiquitous classes of analytical workloads. To process large graph datasets with diverse algorithms, tens of distributed graph processing frameworks emerged. Their users are increasingly expecting high performance for diversifying workloads. Meeting this expectation depends on understanding the performance of each framework. However, performance analysis and characterization of a distributed graph processing framework is challenging. Contributing factors are the irregular nature of graph computation across datasets and algorithms, the semantic gap between workload-level and system-level monitoring, and the lack of lightweight mechanisms for collecting fine-grained performance data. Addressing the challenge, in this work we present Grade10, an experimental framework for fine-grained performance characterization of distributed graph processing workloads. Grade10 captures the graph workload execution as a performance graph from logs and application traces, and builds a fine-grained, unified workload-level and system-level view of performance. Grade10 samples sparsely for lightweight monitoring and addresses the problem of accuracy through a novel approach for resource attribution. Last, it can identify automatically resource bottlenecks and common classes of performance issues. Our real-world experimental evaluation with Giraph and PowerGraph, two state-of-the-art distributed graph processing systems, shows that Grade10 can reveal large differences in the nature and severity of bottlenecks across systems and workloads. We also show that Grade10 can be used in debugging processes, by exemplifying how we find with it a synchronization bug in PowerGraph that slows down affected phases by 1.10-2.50×. Grade10 is an open-source project available at https://github.com/atlarge-research/grade10.
Tim Hegeman, Animesh Trivedi, Alexandru Iosup
CLUSTER3
2020 Is Big Data Performance Reproducible in Modern Cloud Networks?
Alexandru Uta, Alexandru Custura, Dmitry Duplyakin, Ivo Jimenez, Jan S. Rellermeyer, Carlos Maltzahn, Robert Ricci, Alexandru Iosup
NSDI8
2020 3rd Workshop on Hot Topics in Cloud Computing Performance (HotCloudPerf'20): Performance Variability
abstract
No abstract available.
Alexandru Uta, Dmitry Duplyakin, Cristina L. Abad, Nikolas Herbst, Alexandru Iosup
ICPE5
2020 The Workflow Trace Archive: Open-Access Data From Public and Private Computing Infrastructures
abstract
Realistic, relevant, and reproducible experiments often need input traces collected from real-world environments. In this work, we focus on traces of workflows-common in datacenters, clouds, and HPC infrastructures. We show that the state-of-the-art in using workflow-traces raises important issues: (1) the use of realistic traces is infrequent and (2) the use of realistic, open-access traces even more so. Alleviating these issues, we introduce the Workflow Trace Archive (WTA), an open-access archive of workflow traces from diverse computing infrastructures and tooling to parse, validate, and analyze traces. The WTA includes > 48 million workflows captured from > 10 computing infrastructures, representing a broad diversity of trace domains and characteristics. To emphasize the importance of trace diversity, we characterize the WTA contents and analyze in simulation the impact of trace diversity on experiment results. Our results indicate significant differences in characteristics, properties, and workflow structures between workload sources, domains, and fields.
Laurens Versluis, Roland Mathá, Sacheendra Talluri, Tim Hegeman, Radu Prodan, Ewa Deelman, Alexandru Iosup
IEEE Trans. Parallel Distributed Syst.7
2019 The AtLarge Vision on the Design of Distributed Systems and Ecosystems
abstract
High-quality designs of distributed systems and services are essential for our digital economy and society. Threatening to slow down the stream of working designs, we identify the mounting pressure of scale and complexity of (eco-)systems, of ill-defined and wicked problems, and of unclear processes, methods, and tools. We envision design itself as a core research topic in distributed systems, to understand and improve the science and practice of distributed (eco-)system design. Toward this vision, we propose the AtLarge design framework, accompanied by a set of 8 core design principles. We also propose 10 key challenges, which we hope the community can address in the following 5 years. In our experience so far, the proposed framework and principles are practical, and lead to pragmatic and innovative designs for large-scale distributed systems.
Alexandru Iosup, Laurens Versluis, Animesh Trivedi, Erwin Van Eyk, Lucian Toader, Vincent van Beek, Giulia Frascaria, Ahmed Musaafir, Sacheendra Talluri
ICDCS1
2019 Portfolio Scheduling for Managing Operational and Disaster-Recovery Risks in Virtualized Datacenters Hosting Business-Critical Workloads
abstract
Cloud datacenters are increasingly hosting business workloads. Such long-running, on-demand workloads raise important challenges in datacenter operation, requiring efficient online scheduling of workloads with unprecedented characteristics under strict service level agreements (SLAs). In this work, we propose an approach to manage the risk of not meeting SLAs. Our approach is based on portfolio scheduling, which is an online scheduling technique that dynamically selects a scheduling algorithm from a set (portfolio), subject to a possibly changing utility function. Ours is the first datacenter-scheduling approach to consider operational and disaster-recovery risks. Using trace-based simulation with traces collected from a commercial multi-datacenter environment, we give evidence that portfolio scheduling is able to mitigate risks significantly better than its constituent scheduling algorithms and better than datacenter engineers.
Vincent van Beek, Giorgos Oikonomou, Alexandru Iosup
ISPDC3
2019 Graphless: Toward Serverless Graph Processing
abstract
Our society is increasingly solving complex problems through the use of graph processing. Existing graph processing systems focus on performance, which allows addressing ever-larger and more complex problems. They also require uncommon expertise to properly deploy and utilize. To make graph processing generally accessible-to small and medium enterprises and institutions, to common research groups, to individuals-, in this work we design and implement the Graphless graph-processing system. Graphless is based on the serverless paradigm, which proposes to simplify computing by letting developers only focus on small, stateless functions, which are deployed and managed automatically. We address with Graphless the key challenge of combining the stateless functions assumed by serverless computing with the (opposite) data-intensive nature of graph processing. Graphless tackles this challenge through an architectural approach that allows it to deploy with push or with pull operation, and a collection of backend services, such as an orchestrator and a memory-as-a-service component. We implement Graphless and conduct with it real-world experiments using Amazon Lambda for cloud-based serverless resources. Using the LDBC Graphalytics benchmark, we analyze Graphless, and compare its performance and operational cost with the graph-processing systems Apache Giraph (big data domain) and GraphMat (HPC). Overall, we show evidence Graphless provides performance and cost-efficiency similar to Giraph, for algorithms that can benefit from fine-grained elasticity, and lower than GraphMat, but is architecturally easier to deploy, and provides both push and pull operation.
Lucian Toader, Alexandru Uta, Ahmed Musaafir, Alexandru Iosup
ISPDC4
2019 Yardstick: A Benchmark for Minecraft-like Services
abstract
Online gaming applications entertain hundreds of millions of daily active players and often feature vastly complex architecture. Among online games, Minecraft-like games simulate unique (e.g., modifiable) environments, are virally popular, and are increasingly provided as a service. However, the performance of Minecraft-like services, and in particular their scalability, is not well understood. Moreover, currently no benchmark exists for Minecraft-like games. Addressing this knowledge gap, in this work we design and use the Yardstick benchmark to analyze the performance of Minecraft-like services. Yardstick is based on an operational model that captures salient characteristics of Minecraft-like services. As input workload, Yardstick captures important features, such as the most-popular maps used within the Minecraft community. Yardstick captures system- and application-level metrics, and derives from them service-level metrics such as frequency of game-updates under scalable workload. We implement Yardstick, and, through real-world experiments in our clusters, we explore the performance and scalability of popular Minecraft-like servers, including the official vanilla server, and the community-developed servers Spigot and Glowstone. Our findings indicate the scalability limits of these servers, that Minecraft-like services are poorly parallelized, and that Glowstone is the least viable option among those tested.
Jerom van der Sar, Jesse Donkervliet, Alexandru Iosup
ICPE3
2019 Characterization of a Big Data Storage Workload in the Cloud
abstract
The proliferation of big data processing platforms has led to radically different system designs, such as MapReduce and the newer Spark. Understanding the workloads of such systems facilitates tuning and could foster new designs. However, whereas MapReduce workloads have been characterized extensively, relatively little public knowledge exists about the characteristics of Spark workloads in representative environments. To address this problem, in this work we collect and analyze a 6-month Spark workload from a major provider of big data processing services, Databricks. Our analysis focuses on a number of key features, such as the long-term trends of reads and modifications, the statistical properties of reads, and the popularity of clusters and of file formats. Overall, we present numerous findings that could form the basis of new systems studies and designs. Our quantitative evidence and its analysis suggest the existence of daily and weekly load imbalances, of heavy-tailed and bursty behaviour, of the relative rarity of modifications, and of proliferation of big data specific formats.
Sacheendra Talluri, Alicja Luszczak, Cristina L. Abad, Alexandru Iosup
ICPE4
2019 EDITORIAL - Special Issue on Large Scale Cooperative Virtual Environments
Laura Ricci, Alexandru Iosup, Radu Prodan
J. Grid Comput.2
2018 POSUM: A Portfolio Scheduler for MapReduce Workloads
abstract
MapReduce ecosystems are (still) widely popular for big data processing in data centers. To address the diverse non-functional requirements arising from many and increasingly more sophisticated users, the community has developed many scheduling policies for MapReduce workloads. Although some individual policies can dynamically optimize for single and stable performance objectives, such as minimizing runtime or cost, or meeting deadlines for realtime-jobs, it seems unlikely that individual policies will remain competitive for increasingly more dynamic workloads and objectives. In contrast, in this work we investigate the ability to dynamically balance performance and cost of a portfolio scheduler for MapReduce workloads. To this end, we design and implement a portfolio scheduling technique, that is, a system capable of adapting to the current workload characteristics and target objectives by periodically evaluating its set of potential policies, and of switching to "the best" policy that targets the current system state. We implement and evaluate our system with real-world experiments on a workload containing a mixture of real-time and batch jobs, with the purpose of minimizing deadline violations, while keeping batch job slowdown in check. Our results show that POSUM is a promising alternative: it can out-perform the individual policies of its portfolio for the combined optimization goal, even without precise predictions.
Maria A. Voinea, Alexandru Uta, Alexandru Iosup
IEEE BigData3
2018 An Elasticity Study of Distributed Graph Processing
abstract
Graphs are a natural fit for modeling concepts used in solving diverse problems in science, commerce, engineering, and governance. Responding to the variety of graph data and algorithms, many parallel and distributed graph processing systems exist. However, until now these platforms use a static model of deployment: they only run on a pre-defined set of machines. This raises many conceptual and pragmatic issues, including misfit with the highly dynamic nature of graph processing, and could lead to resource waste and high operational costs. In contrast, in this work we explore a dynamic model of deployment. We first characterize workload dynamicity, beyond mere active-vertex variability. Then, to conduct an in-depth elasticity study of distributed graph processing, we build a prototype, JoyGraph, which is the first such system that implements complex, policy-based, and fine-grained elasticity. Using the state-of-the-art LDBC Graphalytics benchmark and the SPEC Cloud Group's elasticity metrics, we show the benefits of elasticity in graph processing: (i) improved resource utilization, (ii) reduced operational costs, and (iii) aligned operation-workload dynamicity. Furthermore, we explore the cost of elasticity in graph processing. We identify a key drawback: although elasticity does not degrade application throughput, graph-processing workloads are sensitive to data movement while leasing or releasing resources.
Sietse Au, Alexandru Uta, Alexey Ilyushkin, Alexandru Iosup
CCGrid4
2018 A Trace-Based Performance Study of Autoscaling Workloads of Workflows in Datacenters
abstract
To improve customer experience, datacenter operators offer support for simplifying application and resource management. For example, running workloads of workflows on behalf of customers is desirable, but requires increasingly more sophisticated autoscaling policies, that is, policies that dynamically provision resources for the customer. Although selecting and tuning autoscaling policies is a challenging task for datacenter operators, so far relatively few studies investigate the performance of autoscaling for workloads of workflows. Complementing previous knowledge, in this work we propose the first comprehensive performance study in the field. Using trace-based simulation, we compare state-of-the-art autoscaling policies across multiple application domains, workload arrival patterns (e.g., burstiness), and system utilization levels. We further investigate the interplay between autoscaling and regular allocation policies, and the complexity cost of autoscaling. Our quantitative study focuses not only on traditional performance metrics and on state-of-the-art elasticity metrics, but also on time-and memory-related autoscaling-complexity metrics. Our main results give strong and quantitative evidence about previously unreported operational behavior, for example, that autoscaling policies perform differently across application domains and allocation and provisioning policies should be co-designed.
Laurens Versluis, Mihai Neacsu, Alexandru Iosup
CCGrid3
2018 Elasticity in Graph Analytics? A Benchmarking Framework for Elastic Graph Processing
abstract
Graphs are a natural fit for modeling concepts used in solving diverse problems in science, commerce, engineering, and governance. Responding to the diversity of graph data and algorithms, many parallel and distributed graph-processing systems exist. However, until now these platforms use a static model of deployment: they only run on a pre-defined set of machines. This raises many conceptual and pragmatic issues, including misfit with the highly dynamic nature of graph processing, and could lead to resource waste and high operational costs. In contrast, in this work we explore the benefits and drawbacks of the dynamic model of deployment. Building a three-layer benchmarking framework for assessing elasticity in graph analytics, we conduct an in-depth elasticity study of distributed graph processing. Our framework is composed of state-of-the-art workloads, autoscalers, and metrics, derived from the LDBC Graphalytics benchmark and SPEC RG Cloud Group's elasticity metrics. We uncover the benefits and cost of elasticity in graph processing: while elasticity allows for fine-grained resource management, and does not degrade application performance, we find that graph workloads are sensitive to data migration while leasing or releasing resources. Moreover, we identify non-trivial interactions between scaling policies and graph workloads, which add an extra level of complexity to resource management and scheduling for graph processing.
Alexandru Uta, Sietse Au, Alexey Ilyushkin, Alexandru Iosup
CLUSTER4
2018 Exploring HPC and Big Data Convergence: A Graph Processing Study on Intel Knights Landing
abstract
The question "Can big data and HPC infrastructure converge?" has important implications for many operators and clients of modern computing. However, answering it is challenging. The hardware is currently different, and fast evolving: big data uses machines with modest numbers of fat cores per socket, large caches, and much memory, whereas HPC uses machines with larger numbers of (thinner) cores, non-trivial NUMA architectures, and fast interconnects. In this work, we investigate the convergence of big data and HPC infrastructure for one of the most challenging application domains, the highly irregular graph processing. We contrast through a systematic, experimental study of over 300,000 core-hours the performance of a modern multicore, Intel Knights Landing (KNL) and of traditional big data hardware, in processing representative graph workloads using state-of-the-art graph analytics platforms. The experimental results indicate KNL is convergence-ready, performance-wise, but only after extensive and expert-level tuning of software and hardware parameters.
Alexandru Uta, Ana Lucia Varbanescu, Ahmed Musaafir, Chris Lemaire, Alexandru Iosup
CLUSTER5
2018 Massivizing Computer Systems: A Vision to Understand, Design, and Engineer Computer Ecosystems Through and Beyond Modern Distributed Systems
abstract
Our society is digital: industry, science, governance, and individuals depend, often transparently, on the inter-operation of large numbers of distributed computer systems. Although the society takes them almost for granted, these computer ecosystems are not available for all, may not be affordable for long, and raise numerous other research challenges. Inspired by these challenges and by our experience with distributed computer systems, we envision Massivizing Computer Systems, a domain of computer science focusing on understanding, controlling, and evolving successfully such ecosystems. Beyond establishing and growing a body of knowledge about computer ecosystems and their constituent systems, the community in this domain should also aim to educate many about design and engineering for this domain, and all people about its principles. This is a call to the entire community: there is much to discover and achieve.
Alexandru Iosup, Alexandru Uta, Laurens Versluis, George Andreadis, Erwin Van Eyk, Tim Hegeman, Sacheendra Talluri, Vincent van Beek, Lucian Toader
ICDCS1
2018 A reference architecture for datacenter scheduling: design, validation, and experiments
George Andreadis, Laurens Versluis, Fabian Mastenbroek, Alexandru Iosup
SC4
2018 A mirroring architecture for sophisticated mobile games using computation-offloading
abstract
Summary Mobile gaming is already a popular and lucrative market. However, the low performance and reduced power capacity of mobile devices severely limit the complexity of mobile games and the duration of their game sessions. To mitigate these issues, in this article, we explore using computation‐offloading, that is, allowing the compute‐intensive parts of mobile games to execute on remote infrastructure. Computation‐offloading raises the combined challenge of addressing the trade‐offs between performance and power‐consumption while also keeping the game playable. We propose Mirror, a system for computation‐offloading that supports the demanding performance requirements of sophisticated mobile games. Mirror proposes several conceptual contributions: support for fine‐grained partitioning, both offline (set by developers) and dynamic (policy‐based), and real‐time asynchronous offloading and user‐input synchronization protocols that enable Mirror‐based systems to bound the delays introduced by offloading and thus to achieve adequate performance. Mirror is compatible with all games that are tick‐based and user‐input deterministic. We implement a real‐world prototype of Mirror and apply it to the real‐world, complex, popular game OpenTTD. The experimental results show that, in comparison with the non‐offloaded OpenTTD, Mirror‐ed OpenTTD can significantly improve performance and power consumption while also delivering smooth gameplay. As a trade‐off, Mirror introduces acceptable delay on user inputs.
M. H. Jiang, Otto W. Visser, I. S. W. B. Prasetya, Alexandru Iosup
Concurr. Comput. Pract. Exp.4
2018 Large Scale Cooperative Virtual Environments
abstract
Large Scale
Laura Ricci, Alexandru Iosup, Radu Prodan
Concurr. Comput. Pract. Exp.2
2018 HPS-HDS: High Performance Scheduling for Heterogeneous Distributed Systems
Florin Pop, Alexandru Iosup, Radu Prodan
Future Gener. Comput. Syst.2
2017 The OpenDC Vision: Towards Collaborative Datacenter Simulation and Exploration for Everybody
abstract
In the new Digital Economy, massive computer systems, often grouped in datacenters, serve as factories "producing" cloud services with massive consumption. However, to afford cloud services globally, we must address new research challenges in designing, operating, and using modern datacenters. We must also address challenges in educating and training the next generation of datacenter engineers. Addressing such challenges, in this work we present our vision on OpenDC: we envision the exploration of various datacenter concepts and technologies, using existing and new scientific methods, enabling new education practices and topics, and leading to the creation of new software and data artifacts. We present the datacenter concepts and technologies we are currently planning to explore using OpenDC. We identify the scientific methods we want to use, and explain our vision of education practices. We present the architecture and open-source program underlying the OpenDC software, and the format and open-access data we use for datacenter experiments. We conclude with an open invitation for the community to join our effort.
Alexandru Iosup, George Andreadis, Vincent van Beek, Matthijs Bijman, Erwin Van Eyk, Mihai Neacsu, Leon Overweel, Sacheendra Talluri, Laurens Versluis, Maaike Visser
ISPDC1
2017 An Experimental Performance Evaluation of Autoscaling Policies for Complex Workflows
abstract
Simplifying the task of resource management and scheduling for customers, while still delivering complex Quality-of-Service (QoS), is key to cloud computing. Many autoscaling policies have been proposed in the past decade to decide on behalf of cloud customers when and how to provision resources to a cloud application utilizing cloud elasticity features. However, in prior work, when a new policy is proposed, it is seldom compared to the state-of-the-art, and is often compared only to static provisioning using a predefined QoS target. This reduces the ability of cloud customers and of cloud operators to choose and deploy an autoscaling policy. In our work, we conduct an experimental performance evaluation of autoscaling policies, using as application model workflows, a commonly used formalism for automating resource management for applications with well-defined yet complex structure. We present a detailed comparative study of general state-of-the-art autoscaling policies, along with two new workflow-specific policies. To understand the performance differences between the 7 policies, we conduct various forms of pairwise and group comparisons. We report both individual and aggregated metrics. Our results highlight the trade-offs between the suggested policies, and thus enable a better understanding of the current state-of-the-art.
Alexey Ilyushkin, Ahmed Ali-Eldin, Nikolas Herbst, Alessandro Vittorio Papadopoulos, Bogdan Ghit, Dick H. J. Epema, Alexandru Iosup
ICPE7
2017 Mirror: A computation-offloading framework for sophisticated mobile games
abstract
The low performance and power limitations of mobile devices severely limit the complexity and the duration of playing sessions of mobile games. This article examines the possibility of using computation-offloading to mitigate these problems while keeping the game playable. We design Mirror, a framework for offloading computation targeted at the demanding performance requirements of sophisticated mobile games. The key conceptual contributions of Mirror are design decisions that allow for dynamic fine-grained client-side offloading decisions, and a protocol for real-time asynchronous offloading for bounding network delays. We implement a prototype of Mirror and test it by performing offloading for the game OpenTTD. The results are promising, showing that Mirror can increase the performance and decrease the power consumption of games while keeping the gameplay fairly smooth.
M. H. Jiang, Otto W. Visser, I. S. W. B. Prasetya, Alexandru Iosup
WoWMoM4
2017 Modeling, analysis, and experimental comparison of streaming graph-partitioning policies
Sungpack Hong, Hassan Chafi, Alexandru Iosup, Dick H. J. Epema
J. Parallel Distributed Comput.4
2016 Design and Experimental Evaluation of Distributed Heterogeneous Graph-Processing Systems
abstract
Graph processing is increasingly used in a variety of domains, from engineering to logistics and from scientific computing to online gaming. To process graphs efficiently, GPU-enabled graph-processing systems such as TOTEM and Medusa exploit the GPU or the combined CPU+GPU capabilities of a single machine. Unlike scalable distributed CPU-based systems such as Pregel and GraphX, existing GPU-enabled systems are restricted to the resources of a single machine, including the limited amount of GPU memory, and thus cannot analyze the increasingly large-scale graphs we see in practice. To address this problem, we design and implement three families of distributed heterogeneous graph-processing systems that can use both the CPUs and GPUs of multiple machines. We further focus on graph partitioning, for which we compare existing graph-partitioning policies and a new policy specifically targeted at heterogeneity. We implement all our distributed heterogeneous systems based on the programming model of the single-machine TOTEM, to which we add (1) a new communication layer for CPUs and GPUs across multiple machines to support distributed graphs, and (2) a workload partitioning method that uses offline profiling to distribute the work on the CPUs and the GPUs. We conduct a comprehensive real-world performance evaluation for all three families. To ensure representative results, we select 3 typical algorithms and 5 datasets with different characteristics. Our results include algorithm run time, performance breakdown, scalability, graph partitioning time, and comparison with other graph-processing systems. They demonstrate the feasibility of distributed heterogeneous graph processing and show evidence of the high performance that can be achieved by combining CPUs and GPUs in a distributed environment.
Ana Lucia Varbanescu, Dick H. J. Epema, Alexandru Iosup
CCGrid4
2016 Which Cloud Auto-Scaler Should I Use for my Application?: Benchmarking Auto-Scaling Algorithms
abstract
Rapid elasticity is one of the essential characteristics of cloud computing identified by NIST [17]. Elasticity allows resources to be provisioned and released to scale rapidly out ward and in ward according to demand. Tens -- if not hundreds -- of algorithms have been proposed in the literature to automatically achieve elastic provisioning [15, 23, 14, 21, 13, 20, 6, 12, 16, 10]. These algorithms are typically referred to as elasticity algorithms, dynamic provisioning techniques or autoscalers. While trying to solve the same problem, sometimes with differing assumption, many of these algorithms are either compared to static provisioning or to a predefined QoS target, e.g., predefined response time target, with very little -- or no -- comparison to previously published work. This reduces the ability of an application owner or a cloud operator to choose and deploy a suitable algorithm from the literature. Many of these algorithms have been tested with one single -- real or synthetic -- workload in a specific use-case [13, 14, 10]. While all published algorithms are shown to work in the specific use-case they were designed for with the, typically short, workloads tested with, it is seldom the case that the real scenarios will be any thing close to the test cases for which the algorithms are shown to work. Bursts occur in workloads occasionally. Workload dynamics change over time and the load-mix of an application significantly affects how provisioning should be done [21].
Ahmed Ali-Eldin, Alexey Ilyushkin, Bogdan Ghit, Nikolas Herbst, Alessandro Vittorio Papadopoulos, Alexandru Iosup
ICPE6
2016 SPEC Research Group's Cloud Working Group: RG Cloud Group
abstract
No abstract available.
Alexandru Iosup, Samuel Kounev, Kai Sachs
ICPE1
2016 Operation analysis of massively multiplayer online games on unreliable resources
Radu Prodan, Alexandru Iosup
Peer-to-Peer Netw. Appl.2
2016 Large scale distributed cooperative environments on clouds and P2P
Laura Ricci, Alexandru Iosup, Radu Prodan
Peer-to-Peer Netw. Appl.2
2016 LDBC Graphalytics: A Benchmark for Large-Scale Graph Analysis on Parallel and Distributed Platforms
abstract
In this paper we introduce LDBC Graphalytics, a new industrial-grade benchmark for graph analysis platforms. It consists of six deterministic algorithms, standard datasets, synthetic dataset generators, and reference output, that enable the objective comparison of graph analysis platforms. Its test harness produces deep metrics that quantify multiple kinds of system scalability, such as horizontal/vertical and weak/strong, and of robustness, such as failures and performance variability. The benchmark comes with open-source software for generating data and monitoring performance. We describe and analyze six implementations of the benchmark (three from the community, three from the industry), providing insights into the strengths and weaknesses of the platforms. Key to our contribution, vendors perform the tuning and benchmarking of their platforms.
Alexandru Iosup, Tim Hegeman, Wing Lung Ngai, Stijn Heldens, Arnau Prat-Pérez, Thomas Manhardt, Hassan Chafi, Mihai Capota, Narayanan Sundaram, Michael J. Anderson, Ilie Gabriel Tanase, Yinglong Xia, Lifeng Nai, Peter Boncz
Proc. VLDB Endow.1
2016 When Game Becomes Life: The Creators and Spectators of Online Game Replays and Live Streaming
abstract
Online gaming franchises such as World of Tanks, Defense of the Ancients, and StarCraft have attracted hundreds of millions of users who, apart from playing the game, also socialize with each other through gaming and viewing gamecasts. As a form of User Generated Content (UGC), gamecasts play an important role in user entertainment and gamer education. They deserve the attention of both industrial partners and the academic communities, corresponding to the large amount of revenue involved and the interesting research problems associated with UGC sites and social networks. Although previous work has put much effort into analyzing general UGC sites such as YouTube, relatively little is known about the gamecast sharing sites. In this work, we provide the first comprehensive study of gamecast sharing sites, including commercial streaming-based sites such as Amazon’s Twitch.tv and community-maintained replay-based sites such as WoTreplays. We collect and share a novel dataset on WoTreplays that includes more than 380,000 game replays, shared by more than 60,000 creators with more than 1.9 million gamers. Together with an earlier published dataset on Twitch.tv, we investigate basic characteristics of gamecast sharing sites, and we analyze the activities of their creators and spectators. Among our results, we find that (i) WoTreplays and Twitch.tv are both fast-consumed repositories, with millions of gamecasts being uploaded, viewed, and soon forgotten; (ii) both the gamecasts and the creators exhibit highly skewed popularity, with a significant heavy tail phenomenon; and (iii) the upload and download preferences of creators and spectators are different: while the creators emphasize their individual skills, the spectators appreciate team-wise tactics. Our findings provide important knowledge for infrastructure and service improvement, for example, in the design of proper resource allocation mechanisms that consider future gamecasting and in the tuning of incentive policies that further help player retention.
Adele Lu Jia, Dick H. J. Epema, Alexandru Iosup
ACM Trans. Multim. Comput. Commun. Appl.4
2015 An Empirical Performance Evaluation of GPU-Enabled Graph-Processing Systems
abstract
Graph processing is increasingly used in knowledge economies and in science, in advanced marketing, social networking, bioinformatics, etc. A number of graph-processing systems, including the GPU-enabled Medusa and Totem, have been developed recently. Understanding their performance is key to system selection, tuning, and improvement. Previous performance evaluation studies have been conducted for CPU-based graph-processing systems, such as Graph and GraphX. Unlike them, the performance of GPU-enabled systems is still not thoroughly evaluated and compared. To address this gap, we propose an empirical method for evaluating GPU-enabled graph-processing systems, which includes new performance metrics and a selection of new datasets and algorithms. By selecting 9 diverse graphs and 3 typical graph-processing algorithms, we conduct a comparative performance study of 3 GPU-enabled systems, Medusa, Totem, and MapGraph. We present the first comprehensive evaluation of GPU-enabled systems with results giving insight into raw processing power, performance breakdown into core components, scalability, and the impact on performance of system-specific optimization techniques and of the GPU generation. We present and discuss many findings that would benefit users and developers interested in GPU acceleration for graph processing.
Ana Lucia Varbanescu, Alexandru Iosup, Dick H. J. Epema
CCGRID3
2015 Statistical Characterization of Business-Critical Workloads Hosted in Cloud Datacenters
abstract
Business-critical workloads -- web servers, mail servers, app servers, etc. -- are increasingly hosted in virtualized data enters acting as Infrastructure-as-a-Service clouds (cloud data enters). Understanding how business-critical workloads demand and use resources is key in capacity sizing, in infrastructure operation and testing, and in application performance management. However, relatively little is currently known about these workloads, because the information is complex -- larges-scale, heterogeneous, shared-clusters -- and because datacenter operators remain reluctant to share such information. Moreover, the few operators that have shared data (e.g., Google and several supercomputing centers) have enabled studies in business intelligence (MapReduce), search, and scientific computing (HPC), but not in business-critical workloads. To alleviate this situation, in this work we conduct a comprehensive study of business-critical workloads hosted in cloud data enters. We collect two large-scale and long-term workload traces corresponding to requested and actually used resources in a distributed datacenter servicing business-critical workloads. We perform an in-depth analysis about workload traces. Our study sheds light into the workload of cloud data enters hosting business-critical workloads. The results of this work can be used as a basis to develop efficient resource management mechanisms for data enters. Moreover, the traces we released in this work can be used for workload verification, modelling and for evaluating resource scheduling policies, etc.
Vincent van Beek, Alexandru Iosup
CCGRID3
2015 An Availability-on-Demand Mechanism for Datacenters
abstract
Data enters are at the core of a wide variety of daily ICT utilities, ranging from scientific computing to online gaming. Due to the scale of today's data enters, the failure of computing resources is a common occurrence that may disrupt the availability of ICT services, leading to revenue loss. Although many high availability (HA) techniques have been proposed to mask resource failures, datacenter users' -- who rent datacenter resources and use them to provide ICT utilities to a global population' -- still have limited management options for dynamically selecting and configuring HA techniques. In this work, we propose Availability-on-Demand (AoD), a mechanism consisting of an API that allows datacenter users to specify availability requirements which can dynamically change, and an availability-aware scheduler that dynamically manages computing resources based on user-specified requirements. The mechanism operates at the level of individual service instance, thus enabling fine-grained control of availability, for example during sudden requirement changes and periodic operations. Through realistic, trace-based simulations, we show that the AoD mechanism can achieve high availability with low cost. The AoD approach consumes about the same CPU hours but with higher availability than approaches which use HA techniques randomly. Moreover, comparing to an ideal approach which has perfect predictions about failures, it consumes 13% to 31% more CPU hours but achieves similar availability for critical parts of applications.
Alexandru Iosup, Assaf Israel, Walfredo Cirne, Danny Raz, Dick H. J. Epema
CCGRID2
2015 Can Portability Improve Performance?: An Empirical Study of Parallel Graph Analytics
abstract
Due to increasingly large datasets, graph analytics - traversals, all-pairs shortest path computations, centrality measures, etc. - are becoming the focus of high-performance computing (HPC). Because HPC is currently dominated by many-core architectures (both CPUs and GPUs), new graph processing solutions have to be defined to efficiently use such computing resources. Prior work focuses on platform-specific performance studies and on platform-specific algorithm development, successfully proving that algorithms highly tuned to GPUs or multi-core CPUs can provide high performance graph analytics. However, the portability of such algorithms remains an important concern for many users, especially the many companies without the resources to invest in HPC or concerned about lock-in in single-use parallel techniques.
Ana Lucia Varbanescu, Merijn Verstraaten, Cees T. A. M. de Laat, Ate Penders, Alexandru Iosup, Henk J. Sips
ICPE5
2015 An Empirical Performance Evaluation of Distributed SQL Query Engines
abstract
Distributed SQL Query Engines (DSQEs) are increasingly used in variety of domains, but especially users in small companies with little expertise may face the challenge of selecting an appropriate engine for their specific applications. Although both industry and academia are attempting to come up with high level benchmarks, the performance of DSQEs has never been explored or compared in-depth. We propose an empirical method for evaluating the performance of DSQEs with representative metrics, datasets, and system configurations. We implement a micro-benchmarking suite of three classes of SQL queries for both a synthetic and a real world dataset and we report response time, resource utilization, and scalability. We use our micro-benchmarking suite to analyze and compare three state-of-the-art engines, viz. Shark, Impala, and Hive. We gain valuable insights for each engine and we present a comprehensive comparison of these DSQEs. We find that different query engines have widely varying performance: Hive is always being outperformed by the other engines, but whether Impala or Shark is the best performer highly depends on the query type.
Stefan van Wouw, José Viña, Alexandru Iosup, Dick H. J. Epema
ICPE3
2015 Socializing by Gaming: Revealing Social Relationships in Multiplayer Online Games
abstract
Multiplayer Online Games (MOGs) like Defense of the Ancients and StarCraft II have attracted hundreds of millions of users who communicate, interact, and socialize with each other through gaming. In MOGs, rich social relationships emerge and can be used to improve gaming services such as match recommendation and game population retention, which are important for the user experience and the commercial value of the companies who run these MOGs. In this work, we focus on understanding social relationships in MOGs. We propose a graph model that is able to capture social relationships of a variety of types and strengths. We apply our model to real-world data collected from three MOGs that contain in total over ten years of behavioral history for millions of players and matches. We compare social relationships in MOGs across different game genres and with regular online social networks like Facebook. Taking match recommendation as an example application of our model, we propose SAMRA, a Socially Aware Match Recommendation Algorithm that takes social relationships into account. We show that our model not only improves the precision of traditional link prediction approaches, but also potentially helps players enjoy games to a higher extent.
Adele Lu Jia, Ruud van de Bovenkamp, Alexandru Iosup, Fernando A. Kuipers, Dick H. J. Epema
ACM Trans. Knowl. Discov. Data4
2015 Area of Simulation: Mechanism and Architecture for Multi-Avatar Virtual Environments
abstract
Although Multi-Avatar Distributed Virtual Environments (MAVEs) such as Real-Time Strategy (RTS) games entertain daily hundreds of millions of online players, their current designs do not scale. For example, even popular RTS games such as the StarCraft series support in a single game instance only up to 16 players and only a few hundreds of avatars loosely controlled by these players, which is a consequence of the Event-Based Lockstep Simulation (EBLS) scalability mechanism they employ. Through empirical analysis, we show that a single Area of Interest (AoI), which is a scalability mechanism that is sufficient for single-avatar virtual environments (such as Role-Playing Games), also cannot meet the scalability demands of MAVEs. To enable scalable MAVEs, in this work we propose Area of Simulation (AoS), a new scalability mechanism, which combines and extends the mechanisms of AoI and EBLS. Unlike traditional AoI approaches, which employ only update-based operational models, our AoS mechanism uses both event-based and update-based operational models to manage not single, but multiple areas of interest. Unlike EBLS, which is traditionally used to synchronize the entire virtual world, our AoS mechanism synchronizes only selected areas of the virtual world. We further design an AoS-based architecture, which is able to use both our AoS and traditional AoI mechanisms simultaneously, dynamically trading-off consistency guarantees for scalability. We implement and deploy this architecture and we demonstrate that it can operate with an order of magnitude more avatars and a larger virtual world without exceeding the resource capacity of players' computers.
Shun-Yun Hu, Alexandru Iosup, Dick H. J. Epema
ACM Trans. Multim. Comput. Commun. Appl.3
2014 V for Vicissitude: The Challenge of Scaling Complex Big Data Workflows
abstract
In this paper we present the scaling of BTWorld, our MapReduce-based approach to observing and analyzing the global BitTorrent network which we have been monitoring for the past 4 years. BTWorld currently provides a comprehensive and complex set of queries implemented in Pig Latin, with data dependencies between them, which translate to several MapReduce jobs that have a heavy-tailed distribution with respect to both execution time and input size characteristics. Processing BitTorrent data in excess of 1 TB with our BTWorld workflow required an in-depth analysis of the entire software stack and the design of a complete optimization cycle. We analyze our system from both theoretical and experimental perspectives and we show how we attained a 15 times larger scale of data processing than our previous results.
Bogdan Ghit, Mihai Capota, Tim Hegeman, Jan Hidders, Dick H. J. Epema, Alexandru Iosup
CCGRID6
2014 KOALA-C: A task allocator for integrated multicluster and multicloud environments
abstract
Companies, scientific communities, and individual scientists with varying requirements for their compute-intensive applications may want to use public Infrastructure-as-a-Service clouds to increase the capacity of the resources they have access to. To enable such access, resource managers that currently act as gateways to clusters may also do so for clouds, but for this they require new architectures and scheduling frameworks. In this paper, we present the design and implementation of KOALA-C, which is an extension of the KOALA multicluster scheduler to multicloud environments. KOALA-C enables uniform management across multicluster and multicloud environments by provisioning resources from both infrastructures and grouping them into clusters of resources called sites. KOALA-C incorporates a comprehensive list of policies for scheduling jobs across multiple (sets of) sites, including both traditional policies and two new policies inspired by the well-known TAGS task assignment policy in distributed-server systems. Finally, we evaluate KOALA-C through realistic simulations and real-world experiments, and show that the new architecture and in particular its new policies show promise in achieving good job slowdown with high resource utilization.
Lipu Fei, Bogdan Ghit, Alexandru Iosup, Dick H. J. Epema
CLUSTER3
2014 How Well Do Graph-Processing Platforms Perform? An Empirical Performance Evaluation and Analysis
abstract
Graph-processing platforms are increasingly used in a variety of domains. Although both industry and academia are developing and tuning graph-processing algorithms and platforms, the performance of graph-processing platforms has never been explored or compared in-depth. Thus, users face the daunting challenge of selecting an appropriate platform for their specific application. To alleviate this challenge, we propose an empirical method for benchmarking graph-processing platforms. We define a comprehensive process, and a selection of representative metrics, datasets, and algorithmic classes. We implement a benchmarking suite of five classes of algorithms and seven diverse graphs. Our suite reports on basic (user-lever) performance, resource utilization, scalability, and various overhead. We use our benchmarking suite to analyze and compare six platforms. We gain valuable insights for each platform and present the first comprehensive comparison of graph-processing platforms.
Marcin Biczak, Ana Lucia Varbanescu, Alexandru Iosup, Claudio Martella, Theodore L. Willke
IPDPS4
2014 Dynamic Resource Management in Cloud-based Distributed Virtual Environments
abstract
As an elastic hosting platform, cloud computing has been attracting many attentions for transferring compute-intensive applications from static self-hosting to flexible cloud-based hosting. Distributed virtual environments (DVEs) which typically involve massive users interacting at the same time and feature significant workload dynamics either in spatial due to in-game user mobility or in temporal due to the fluctuating user population, potentially are suitable applications with cloud-based hosting because of the need of resource elasticity. We explore the dynamic resource management for cloud-based DVEs by taking into account their multi-level workload dynamics which differ them from other applications. Simulation results demonstrates the advantages of our developed methods over existing ones.
Yunhua Deng, Zhe Huang 0004, Alexandru Iosup, Rynson W. H. Lau
ACM Multimedia4
2014 Characterization of Human Mobility in Networked Virtual Environments
abstract
The design and tuning of networked virtual environments (NVEs), such as World of Warcraft (WoW), require understanding the in-NVE mobility characteristics of their citizens. Although many mobility-aware NVE systems already exist, their validation and further development have been hampered by the lack of public datasets and of comparison studies based on multiple datasets. To address these two issues, in this work we collect from WoW mobility traces for over 30,000 virtual citizens, and compare these traces with traces collected from Second Life (SL) where the environment is designed and changed significantly by the citizens themselves. Furthermore, motivated by the existence of numerous studies and models of networked real-world environments (NRE), we systematically compare the characteristics of two NVE and two NRE mobility traces. Our comparative study reveals that long-tail distributions characterize well various mobility characteristics, that the invisible boundary of human movement also appears for NVEs, and that area-visitation shows personal preferences. We also find several differences between NVE and NRE mobility characteristics.
Niels Brouwers, Alexandru Iosup, Dick H. J. Epema
NOSSDAV3
2014 An experience report on using gamification in technical higher education
abstract
Technical universities, especially in Europe, are facing an important challenge in attracting more diverse groups of students, and in keeping the students they attract motivated and engaged in the curriculum. We describe our experience with gamification, which we loosely define as a teaching technique that uses social gaming elements to deliver higher education. Over the past three years, we have applied gamification to undergraduate and graduate courses in a leading technical university in the Netherlands and in Europe. Ours is one of the first long-running attempts to show that gamification can be used to teach technically challenging courses. The two gamification-based courses, the first-year B.Sc. course Computer Organization and an M.Sc.-level course on the emerging technology of Cloud Computing, have been cumulatively followed by over 450 students and passed by over 75% of them, at the first attempt. We find that gamification is correlated with an increase in the percentage of passing students, and in the participation in voluntary activities and challenging assignments. Gamification seems to also foster interaction in the classroom and trigger students to pay more attention to the design of the course. We also observe very positive student assessments and volunteered testimonials, and a Teacher of the Year award.
Alexandru Iosup, Dick H. J. Epema
SIGCSE1
2014 Balanced resource allocations across multiple dynamic MapReduce clusters
abstract
Running multiple instances of the MapReduce framework concurrently in a multicluster system or datacenter enables data, failure, and version isolation, which is attractive for many organizations. It may also provide some form of performance isolation, but in order to achieve this in the face of time-varying workloads submitted to the MapReduce instances, a mechanism for dynamic resource (re-)allocations to those instances is required. In this paper, we present such a mechanism called Fawkes that attempts to balance the allocations to MapReduce instances so that they experience similar service levels. Fawkes proposes a new abstraction for deploying MapReduce instances on physical resources, the MR-cluster, which represents a set of resources that can grow and shrink, and that has a core on which MapReduce is installed with the usual data locality assumptions but that relaxes those assumptions for nodes outside the core. Fawkes dynamically grows and shrinks the active MR-clusters based on a family of weighting policies with weights derived from monitoring their operation.
Bogdan Ghit, Nezih Yigitbasi, Alexandru Iosup, Dick H. J. Epema
SIGMETRICS3
2014 Benchmarking graph-processing platforms: a vision
abstract
Processing graphs, especially at large scale, is an increasingly useful activity in a variety of business, engineering, and scientific domains. Already, there are tens of graph-processing platforms, such as Hadoop, Giraph, GraphLab, etc., each with a different design and functionality. For graph-processing to continue to evolve, users have to find it easy to select a graph-processing platform, and developers and system integrators have to find it easy to quantify the performance and other non-functional aspects of interest. However, the state of performance analysis of graph-processing platforms is still immature: there are few studies and, for the few that exist, there are few similarities, and relatively little understanding of the impact of dataset and algorithm diversity on performance. Our vision is to develop, with the help of the performance-savvy community, a comprehensive benchmarking suite for graph-processing platforms. In this work, we take a step in this direction, by proposing a set of seven challenges, summarizing our previous work on performance evaluation of distributed graph-processing platforms, and introducing our on-going work within the SPEC Research Group's Cloud Working Group.
Ana Lucia Varbanescu, Alexandru Iosup, Claudio Martella, Theodore L. Willke
ICPE3
2014 SLA-based operations of massively multiplayer online games in clouds
Vlad Nae, Radu Prodan, Alexandru Iosup
Multim. Syst.3
2013 The BTWorld use case for big data analytics: Description, MapReduce logical workflow, and empirical evaluation
abstract
The commoditization of big data analytics, that is, the deployment, tuning, and future development of big data processing platforms such as MapReduce, relies on a thorough understanding of relevant use cases and workloads. In this work we propose BTWorld, a use case for time-based big data analytics that is representative for processing data collected periodically from a global-scale distributed system. BTWorld enables a data-driven approach to understanding the evolution of BitTorrent, a global file-sharing network that has over 100 million users and accounts for a third of today's upstream traffic. We describe for this use case the analyst questions and the structure of a multi-terabyte data set. We design a MapReduce-based logical workflow, which includes three levels of data dependency - inter-query, inter-job, and intra-job - and a query diversity that make the BTWorld use case challenging for today's big data processing tools; the workflow can be instantiated in various ways in the MapReduce stack. Last, we instantiate this complex workflow using Pig-Hadoop-HDFS and evaluate the use case empirically. Our MapReduce use case has challenging features: small (kilobytes) to large (250 MB) data sizes per observed item, excellent (10-6) and very poor (102) selectivity, and short (seconds) to long (hours) job duration.
Tim Hegeman, Bogdan Ghit, Mihai Capota, Jan Hidders, Dick H. J. Epema, Alexandru Iosup
IEEE BigData6
2013 Towards an Optimized Big Data Processing System
abstract
Scalable by design to very large computing systems such as grids and clouds, MapReduce is currently a major big data processing paradigm. Nevertheless, existing performance models for MapReduce only comply with specific workloads that process a small fraction of the entire data set, thus failing to assess the capabilities of the MapReduce paradigm under heavy workloads that process exponentially increasing data volumes. The goal of my PhD is to build and analyze a scalable and dynamic big data processing system, including storage (distributed file system), execution engine (MapReduce), and query language (Pig). My contributions for the first two years of PhD research are the following: 1) the design and implementation of a resource management system part of a MapReduce-based processing system for deploying and resizing MapReduce clusters over multicluster systems, 2) the design and implementation of a benchmarking tool for the MapReduce processing system, and 3) the evaluation and modeling of MapReduce using workloads with very large data sets. Furthermore, based on the first two years research, we will optimize the MapReduce system to efficiently process terabytes of data.
Bogdan Ghit, Alexandru Iosup, Dick H. J. Epema
CCGRID2
2013 Extending the Capabilities of Mobile Devices for Online Social Applications through Cloud Offloading
abstract
Handheld devices are becoming an attractive option for users to interact with their social network, through online social applications. We are witnessing a rapid adoption of smarter devices all around us, which brings with it orders of magnitude in heterogeneity. Thus, researchers in the field of distributed systems are faced with new challenges: How to optimize performance for devices that are so diverse in terms of energy consumption, processing power and communication capabilities? My PhD research focuses on this challenge, adopting techniques for offloading operations from mobile to more powerful cloud-based infrastructure, and brings a three-fold contribution. First, we have characterized and modeled workloads of online social applications, and empirically validated them using traces of hundreds of real applications. Second, we are currently investigating offloading mechanisms, including: communication offloading, lossy performance offloading, and loss less performance offloading. We have been testing and evaluating these mechanisms with several mobile applications, measuring performance and energy consumption. Third, we will create an integrated cloud based offloading system that aims to improve the performance of online social applications. We will empirically evaluate this system using both simulations and open-source real-world applications.
Alexandru-Corneliu Olteanu, Nicolae Tapus, Alexandru Iosup
CCGRID3
2013 Massivizing Multi-player Online Games on Clouds
abstract
Massively Multiplayer Online Games (MMOGs) are an important type of distributed applications and have millions of users. Traditionally, MMOGs are hosted on dedicated clusters, distributed globally. With the advent of cloud computing, MMOGs such as Zynga's are increasingly run on cloud resources, through the use of cloud technology and innovation. Massivizing MMOGs on clouds is the focus of my PhD research. My main contributions are: 1) analyzing and modeling various MMOG workloads, including those of social and traditional real-time games, 2) designing and implementing a cost-efficient and reliable cloud-based MMOG platform, 3) designing and implementing a scalable MMOG system which employs domain-specific scaling techniques to support the real time strategy games of the future, 4) experimental prototypes and tools to evaluate our proposed research via real-world experimentation and simulation, and applying our proposed research to a popular real-world application. In this article, I introduce my research progress and my future plans.
Alexandru Iosup, Dick H. J. Epema
CCGRID2
2013 Scheduling Jobs in the Cloud Using On-Demand and Reserved Instances
Kefeng Deng, Alexandru Iosup, Dick H. J. Epema
Euro-Par3
2013 A Periodic Portfolio Scheduler for Scientific Computing in the Data Center
Kefeng Deng, Ruben Verboon, Kaijun Ren, Alexandru Iosup
JSSPP4
2013 Exploring portfolio scheduling for long-term execution of scientific workloads in IaaS clouds
abstract
Long-term execution of scientific applications often leads to dynamic workloads and varying application requirements. When the execution uses resources provisioned from IaaS clouds, and thus consumption-related payment, efficient and online scheduling algorithms must be found. Portfolio scheduling, which selects dynamically a suitable policy from a broad portfolio, may provide a solution to this problem. However, selecting online the right policy from possibly tens of alternatives remains challenging. In this work, we introduce an abstract model to explore this selection problem. Based on the model, we present a comprehensive portfolio scheduler that includes tens of provisioning and allocation policies. We propose an algorithm that can enlarge the chance of selecting the best policy in limited time, possibly online. Through trace-based simulation, we evaluate various aspects of our portfolio scheduler, and find performance improvements from 7% to 100% in comparison with the best constituent policies and high improvement for bursty workloads.
Kefeng Deng, Junqiang Song, Kaijun Ren, Alexandru Iosup
SC4
2013 Towards a workload model for online social applications: ICPE 2013 work-in-progress paper
abstract
Popular online social applications hosted by social platforms serve, each, millions of interconnected users. Understanding the workloads of these applications is key in improving the management of their performance and costs. In this work, we analyse traces gathered over a period of thirty-one months for hundreds of Facebook applications. We characterize the popularity of applications, which describes how applications attract users, and the evolution pattern, which describes how the number of users changes over the lifetime of an application. We further model both application popularity and evolution, and validate our model statistically, by fitting five probability distributions to empirical data for each of the model variables. Among the results, we find that most applications reach their maximum number of users within a third of their lifetime, and that the lognormal distribution provides the best fit for the popularity distribution.
Alexandru-Corneliu Olteanu, Alexandru Iosup, Nicolae Tapus
ICPE2
2013 The Failure Trace Archive: Enabling the comparison of failure measurements and models of distributed systems
Bahman Javadi, Derrick Kondo, Alexandru Iosup, Dick H. J. Epema
J. Parallel Distributed Comput.3
2013 Procedural content generation for games: A survey
abstract
Hundreds of millions of people play computer games every day. For them, game content—from 3D objects to abstract puzzles—plays a major entertainment role. Manual labor has so far ensured that the quality and quantity of game content matched the demands of the playing community, but is facing new scalability challenges due to the exponential growth over the last decade of both the gamer population and the production costs. Procedural Content Generation for Games (PCG-G) may address these challenges by automating, or aiding in, game content generation. PCG-G is difficult, since the generator has to create the content, satisfy constraints imposed by the artist, and return interesting instances for gamers. Despite a large body of research focusing on PCG-G, particularly over the past decade, ours is the first comprehensive survey of the field of PCG-G. We first introduce a comprehensive, six-layered taxonomy of game content: bits, space, systems, scenarios, design, and derived. Second, we survey the methods used across the whole field of PCG-G from a large research body. Third, we map PCG-G methods to game content layers; it turns out that many of the methods used to generate game content from one layer can be used to generate content from another. We also survey the use of methods in practice, that is, in commercial or prototype games. Fourth and last, we discuss several directions for future research in PCG-G, which we believe deserve close attention in the near future.
Mark Hendrikx, Sebastiaan A. Meijer, Joeri Van Der Velden, Alexandru Iosup
ACM Trans. Multim. Comput. Commun. Appl.4
2012 An Analysis of Provisioning and Allocation Policies for Infrastructure-as-a-Service Clouds
abstract
Today, many commercial and private cloud computing providers offer resources for leasing under the infrastructure as a service (IaaS) paradigm. Although an abundance of mechanisms already facilitate the lease and use of single infrastructure resources, to complete multi-job workloads IaaS users still need to select adequate provisioning and allocation policies to instantiate resources and map computational jobs to them. While such policies have been studied in the past, no experimental investigation in the context of clouds currently exists that considers them jointly. In this paper we present a comprehensive and empirical performance-cost analysis of provisioning and allocation policies in IaaS clouds. We first introduce a taxonomy of both types of policies, based on the type of information used in the decision process, and map to this taxonomy eight provisioning and four allocation policies. Then, we analyze the performance and cost of these policies through experimentation in three clouds, including Amazon EC2. We show that policies that dynamically provision and/or allocate resources can achieve better performance and cost. Finally, we also look at the interplay between provisioning and allocation, for which we show preliminary results.
David Villegas, Athanasios Antoniou, Seyed Masoud Sadjadi, Alexandru Iosup
CCGRID4
2012 ExPERT: Pareto-Efficient Task Replication on Grids and a Cloud
abstract
Many scientists perform extensive computations by executing large bags of similar tasks (BoTs) in mixtures of computational environments, such as grids and clouds. Although the reliability and cost may vary considerably across these environments, no tool exists to assist scientists in the selection of environments that can both fulfill deadlines and fit budgets. To address this situation, we introduce the Expert BoT scheduling framework. Our framework systematically selects from a large search space the Pareto-efficient scheduling strategies, that is, the strategies that deliver the best results for both make span and cost. Expert chooses from them the best strategy according to a general, user-specified utility function. Through simulations and experiments in real production environments, we demonstrate that Expert can substantially reduce both make span and cost in comparison to common scheduling strategies. For bioinformatics BoTs executed in a real mixed grid + cloud environment, we show how the scheduling strategy selected by Expert reduces both make span and cost by 30%-70%, in comparison to commonly-used scheduling strategies.
Orna Agmon Ben-Yehuda, Assaf Schuster, Artyom Sharov, Mark Silberstein, Alexandru Iosup
IPDPS5
2011 On the Performance Variability of Production Cloud Services
abstract
Cloud computing is an emerging infrastructure paradigm that promises to eliminate the need for companies to maintain expensive computing hardware. Through the use of virtualization and resource time-sharing, clouds address with a single set of physical resources a large user base with diverse needs. Thus, clouds have the potential to provide their owners the benefits of an economy of scale and, at the same time, become an alternative for both the industry and the scientific community to self-owned clusters, grids, and parallel production environments. For this potential to become reality, the first generation of commercial clouds need to be proven to be dependable. In this work we analyze the dependability of cloud services. Towards this end, we analyze long-term performance traces from Amazon Web Services and Google App Engine, currently two of the largest commercial clouds in production. We find that the performance of about half of the cloud services we investigate exhibits yearly and daily patterns, but also that most services have periods of especially stable performance. Last, through trace-based simulation we assess the impact of the variability observed for the studied cloud services on three large-scale applications, job execution in scientific computing, virtual goods trading in social networks, and state management in social gaming. We show that the impact of performance variability depends on the application, and give evidence that performance variability can be an important factor in cloud provider selection.
Alexandru Iosup, Nezih Yigitbasi, Dick H. J. Epema
CCGRID1
2011 Identifying, analyzing, and modeling flashcrowds in BitTorrent
abstract
Flashcrowds - sudden surges of user arrivals - do occur in BitTorrent, and they can lead to severe service deprivation. However, very little is known about their occurrence patterns and their characteristics in real-world deployments, and many basic questions about BitTorrent flashcrowds, such as How often do they occur? and How long do they last?, remain unanswered. In this paper, we address these questions by studying three datasets that cover millions of swarms from two of the largest BitTorrent trackers. We first propose a model for BitTorrent flashcrowds and a procedure for identifying, analyzing, and modeling BitTorrent flashcrowds. Then we evaluate quantitatively the impact of flashcrowds on BitTorrent users, and we develop an algorithm that identifies BitTorrent flashcrowds. Finally, we study statistically the properties of BitTorrent flashcrowds identified from our datasets, such as their arrival time, duration, and magnitude, and we investigate the relationship between flashcrowds and swarm growth, and the arrival rate of flashcrowds in BitTorrent trackers. In particular, we find that BitTorrent flashcrowds only occur in very small fractions (0.3-2%) of the swarms but that they can affect over ten million users.
Boxun Zhang, Alexandru Iosup, Johan A. Pouwelse, Dick H. J. Epema
Peer-to-Peer Computing2
2011 A new business model for massively multiplayer online games
abstract
Today, highly successful Massively Multiplayer Online Games (MMOGs) have millions of registered users and hundreds of thousands of active concurrent players. To sustain their highly variable load, game operators over-provision a large static infrastructure capable of sustaining the game peak load, even though a large portion of the resources is unused most of the time. This inefficient resource utilisation has negative economic impacts by preventing any but the largest hosting centres from joining the market and dramatically increases prices.In this paper, we propose a new business model of hosting and operating MMOGs based on Cloud computing principles involving four actors: resource provider, game operator, game provider, and client. Our model efficiently provisions on-demand virtualised resources to game sessions based on their dynamic client load, which dramatically decreases prices and gives small and medium enterprises the opportunity of joining the market through zero initial investment.We validate our new model and its underlying business relationships through trace-based simulations utilising six months worth of monitoring data from a real-life MMOG using emulated resources from 16 of the largest Cloud resource providers currently on the market. We demonstrate that our model can operate state-of-the-art MMOGs with an average monthly gross profit of nearly $6 million excluding game purchase prices, overheads and taxation, while being able to maintain and control the QoS offered to all clients. Finally, we show how our approach is capable of operating next generation very highly interactive MMOGs with a small increase of 5.8% in the subscription price.
Vlad Nae, Radu Prodan, Alexandru Iosup, Thomas Fahringer
ICPE3
2011 POGGI: generating puzzle instances for online games on grid infrastructures
abstract
Abstract The developers of Massively multiplayer online games (MMOGs) rely on attractive content such as logical challenges (puzzles) to entertain and generate revenue from millions of concurrent players. While a large part of the content is currently generated by hand, the exponentially increasing number of players and their new demand for player‐customized content have made manual content generation undesirable. Thus, in this work we investigate the problem of automated, player‐customized game content generation at MMOG scale; to make our investigation tractable, we focus on puzzle game content generation. In this context, we design POGGI, an architecture for generating player‐customized game content using on‐demand, grid resources. POGGI focuses on the large‐scale generation of puzzle instances that match well the solving ability of the players, and that lead to fresh playing experience. Using our reference implementation of POGGI, we show through experiments in a real resource pool of over 1600 nodes that POGGI can generate commercial‐quality content at MMOG scale. Copyright © 2010 John Wiley & Sons, Ltd.
Alexandru Iosup
Concurr. Comput. Pract. Exp.1
2011 Performance Analysis of Cloud Computing Services for Many-Tasks Scientific Computing
abstract
Cloud computing is an emerging commercial infrastructure paradigm that promises to eliminate the need for maintaining expensive computing facilities by companies and institutes alike. Through the use of virtualization and resource time sharing, clouds serve with a single set of physical resources a large user base with different needs. Thus, clouds have the potential to provide to their owners the benefits of an economy of scale and, at the same time, become an alternative for scientists to clusters, grids, and parallel production environments. However, the current commercial clouds have been built to support web and small database workloads, which are very different from typical scientific computing workloads. Moreover, the use of virtualization and resource time sharing may introduce significant performance penalties for the demanding scientific computing workloads. In this work, we analyze the performance of cloud computing services for scientific computing workloads. We quantify the presence in real scientific computing workloads of Many-Task Computing (MTC) users, that is, of users who employ loosely coupled applications comprising many tasks to achieve their scientific goals. Then, we perform an empirical evaluation of the performance of four commercial cloud computing services including Amazon EC2, which is currently the largest commercial cloud. Last, we compare through trace-based simulation the performance characteristics and cost models of clouds and other scientific computing platforms, for general and MTC-based scientific computing workloads. Our results indicate that the current clouds need an order of magnitude in performance improvement to be useful to the scientific community, and show which improvements should be considered first to address this discrepancy between offer and demand.
Alexandru Iosup, Simon Ostermann 0001, Nezih Yigitbasi, Radu Prodan, Thomas Fahringer, Dick H. J. Epema
IEEE Trans. Parallel Distributed Syst.1
2011 Dynamic Resource Provisioning in Massively Multiplayer Online Games
abstract
Today's Massively Multiplayer Online Games (MMOGs) can include millions of concurrent players spread across the world and interacting with each other within a single session. Faced with high resource demand variability and with misfit resource renting policies, the current industry practice is to overprovision for each game tens of self-owned data centers, making the market entry affordable only for big companies. Focusing on the reduction of entry and operational costs, we investigate a new dynamic resource provisioning method for MMOG operation using external data centers as low-cost resource providers. First, we identify in the various types of player interaction a source of short-term load variability, which complements the long-term load variability due to the size of the player population. Then, we introduce a combined MMOG processor, network, and memory load model that takes into account both the player interaction type and the population size. Our model is best used for estimating the MMOG resource demand dynamically, and thus, for dynamic resource provisioning based on the game world entity distribution. We evaluate several classes of online predictors for MMOG entity distribution and propose and tune a neural network-based predictor to deliver good accuracy consistently under real-time performance constraints. We assess using trace-based simulation the impact of the data center policies on the quality of resource provisioning. We find that the dynamic resource provisioning can be much more efficient than its static alternative even when the external data centers are busy, and that data centers with policies unsuitable for MMOGs are penalized by our dynamic resource provisioning method. Finally, we present experimental results showing the real-time parallelization and load balancing of a real game prototype using data center resources provisioned using our method and show its advantage against a rudimentary client threshold approach.
Vlad Nae, Alexandru Iosup, Radu Prodan
IEEE Trans. Parallel Distributed Syst.2
2010 The Failure Trace Archive: Enabling Comparative Analysis of Failures in Diverse Distributed Systems
abstract
With the increasing functionality and complexity of distributed systems, resource failures are inevitable. While numerous models and algorithms for dealing with failures exist, the lack of public trace data sets and tools has prevented meaningful comparisons. To facilitate the design, validation, and comparison of fault-tolerant models and algorithms, we have created the Failure Trace Archive (FTA) as an online public repository of availability traces taken from diverse parallel and distributed systems. Our main contributions in this study are the following. First, we describe the design of the archive, in particular the rationale of the standard FTA format, and the design of a toolbox that facilitates automated analysis of trace data sets. Second, applying the toolbox, we present a uniform comparative analysis with statistics and models of failures in nine distributed systems. Third, we show how different interpretations of these data sets can result in different conclusions. This emphasizes the critical need for the public availability of trace data and methods for their analysis.
Derrick Kondo, Bahman Javadi, Alexandru Iosup, Dick H. J. Epema
CCGRID3
2010 A Model for Space-Correlated Failures in Large-Scale Distributed Systems
Matthieu Gallet, Nezih Yigitbasi, Bahman Javadi, Derrick Kondo, Alexandru Iosup, Dick H. J. Epema
Euro-Par (1)5
2010 Sampling Bias in BitTorrent Measurements
Boxun Zhang, Alexandru Iosup, Johan A. Pouwelse, Dick H. J. Epema, Henk J. Sips
Euro-Par (1)2
2010 Performance analysis of dynamic workflow scheduling in multicluster grids
abstract
Scientists increasingly rely on the execution of workflows in grids to obtain results from complex mixtures of applications. However, the inherently dynamic nature of grid workflow scheduling, stemming from the unavailability of scheduling information and from resource contention among the (multiple) workflows and the non-workflow system load, may lead to poor or unpredictable performance. In this paper we present a comprehensive and realistic investigation of the performance of a wide range of dynamic workflow scheduling policies in multicluster grids. We first introduce a taxonomy of grid workflow scheduling policies that is based on the amount of dynamic information used in the scheduling process, and map to this taxonomy seven such policies across the full spectrum of information use. Then, we analyze the performance of these scheduling policies through simulations and experiments in a real multicluster grid. We find that there is no single grid workflow scheduling policy with good performance across all the investigated scenarios. We also find from our real system experiments that with demanding workloads, the limitations of the head-nodes of the grid clusters may lead to performance loss not expected from the simulation results. We show that task throttling, that is, limiting the per-workflow number of tasks dispatched to the system, prevents the head-nodes from becoming overloaded while largely preserving performance, at least for communication-intensive workflows. 1.
Omer Ozan Sonmez, Nezih Yigitbasi, Saeid Abrishami, Alexandru Iosup, Dick H. J. Epema
HPDC4
2010 BTWorld: towards observing the global BitTorrent file-sharing network
abstract
Today, the BitTorrent Peer-to-Peer file-sharing network is one of the largest Internet applications---it generates massive traffic volumes, it is deployed in thousands of independent communities, and it serves millions of unique users worldwide. Despite a large number of empirical and theoretical studies, observing the state of the global BitTorrent network remains a grand challenge for the BitTorrent community. To address this challenge, in this work we introduce BT-World, an architecture for observing the global BitTorrent network without help from the ISPs. We design BTWorld around three main features specific to BitTorrent measurements. First, our architecture is able to find public trackers, that is, the BitTorrent components that offer unrestricted service to peers around the world. Second, by observing the state of these trackers, BTWorld obtains information about the performance, scalability, and reliability of BitTorrent. Third, BTWorld is designed to pre-process the large volumes of recorded data for later analysis. We demonstrate the viability of our architecture by deploying it in practice, to observe and analyze one week of operation of a large part of the global BitTorrent network--over 10 million swarms and tens of millions of concurrent users. We also show that BT-World can shed light on BitTorrent phenomena, such as the presence of spam trackers and giant swarms.
Maciej Wojciechowski, Mihai Capota, Johan A. Pouwelse, Alexandru Iosup
HPDC4
2009 Scheduling Strategies for Cycle Scavenging in Multicluster Grid Systems
abstract
The use of today's multicluster grids exhibits periods of submission bursts with periods of normal use and even of idleness. To avoid resource contention, many users employ observational scheduling, that is, they postpone the submission of relatively low-priority jobs until a cluster becomes (largely) idle. However, observational scheduling leads to resource contention when several such users crowd the same idle cluster. Moreover, this job execution model either delays the execution of more important jobs, or requires extensive administrative support for job and user priorities. Instead, in this work we investigate the use of cycle scavenging to run jobs on grid resources politely yet efficiently, and with an acceptable administrative cost. We design a two-level cycle scavenging scheduling architecture that runs unobtrusively alongside regular grid scheduling. We equip this scheduler with two novel cycle scavenging scheduling policies that enforce fair resource sharing among competing cycle scavenging users. We show through experiments with real and synthetic applications in a real multicluster grid that the proposed architecture can execute jobs politely yet efficiently.
Omer Ozan Sonmez, Bart Grundeken, Hashim H. Mohamed, Alexandru Iosup, Dick H. J. Epema
CCGRID4
2009 C-Meter: A Framework for Performance Analysis of Computing Clouds
abstract
Cloud computing has emerged as a new technology that provides large amounts of computing and data storage capacity to its users with a promise of increased scalability, high availability, and reduced administration and maintenance costs. As the use of cloud computing environments increases, it becomes crucial to understand the performance of these environments. So, it is of great importance to assess the performance of computing clouds in terms of various metrics, such as the overhead of acquiring and releasing the virtual computing resources, and other virtualization and network communications overheads. To address these issues, we have designed and implemented C-Meter, which is a portable, extensible, and easy-to-use framework for generating and submitting test workloads to computing clouds. In this paper, first we state the requirements for frameworks to assess the performance of computing clouds. Then, we present the architecture of the C-Meter framework and discuss several cloud resource management alternatives. Finally, we present our early experiences with C-Meter in Amazon EC2. We show how C-Meter can be used for assessing the overhead of acquiring and releasing the virtual computing resources, for comparing different configurations, and for evaluating different scheduling algorithms.
Nezih Yigitbasi, Alexandru Iosup, Dick H. J. Epema, Simon Ostermann 0001
CCGRID2
2009 An Aspect-Oriented Approach for Disaster Prevention Simulation Workflows on Supercomputers, Clusters, and Grids
abstract
Computer simulation is an important factor in today's disaster prevention procedures. Simulation codes assess the evolution and impact of various physical phenomena in domains such as nuclear and environmental sciences, and ultimately help saving lives. However, new and more computationally demanding models, and new regulations for personnel training have increased the demand for computational power. While the existing simulation codes can be ported to computing environments that can meet the new demand, such as supercomputers, clusters, and grids, it is too expensive and time-consuming to rewrite and re-certify them. Instead, in this work we propose an aspect-oriented approach that takes existing simulation functionality and combines it with functionality required to run the simulation on different computing environments transparently to the simulation developer. Through experiments in the DAS-3 multi-cluster grid we show that our approach increases the reusability, the maintainability, the scalability, and the robustness of a real disaster prevention simulation, while incurring a low performance overhead.
Tudor B. Ionescu, Andreas Piater, Walter Scheuermann, Eckart Laurien, Alexandru Iosup
DS-RT5
2009 Introduction
Thomas Fahringer, Alexandru Iosup, Marian Bubak, Matei Ripeanu, Xian-He Sun, Hong Linh Truong 0001
Euro-Par2
2009 POGGI: Puzzle-Based Online Games on Grid Infrastructures
Alexandru Iosup
Euro-Par1
2009 Trace-based evaluation of job runtime and queue wait time predictions in grids
abstract
Large-scale distributed computing systems such as grids are serving a growing number of scientists. These environments bring about not only the advantages of an economy of scale, but also the challenges of resource and workload heterogeneity. A consequence of these two forms of heterogeneity is that job runtimes and queue wait times are highly variable, which generally reduces system performance and makes grids difficult to use by the common scientist. Predicting job runtimes and queue wait times have been widely studied for parallel environments. However, there is no detailed investigation on how the proposed prediction methods perform in grids, whose resource structure and workload characteristics are very different from those in parallel systems. In this paper, we assess the performance and benefit of predicting job runtimes and queue wait times in grids based on traces gathered from various research and production grid environments. First, we evaluate the performance of simple yet widely used time series prediction methods and the effect of applying them to different types of job classes (e.g., all jobs submitted by single users or to single sites). Then, we investigate the performance of two kinds of queue wait time prediction methods for grids. Last, we investigate whether prediction-based grid-level scheduling policies can have better performance than policies that do not use predictions.
Omer Ozan Sonmez, Nezih Yigitbasi, Alexandru Iosup, Dick H. J. Epema
HPDC3
2008 DGSim: Comparing Grid Resource Management Architectures through Trace-Based Simulation
Alexandru Iosup, Omer Ozan Sonmez, Dick H. J. Epema
Euro-Par1
2008 The performance of bags-of-tasks in large-scale distributed systems
abstract
Ever more scientists are employing large-scale distributed systems such as grids for their computational work, instead of tightly coupled high-performance computing systems. However, while these distributed systems are more cost-effective, their heterogeneity in terms of hardware, software, and systems administration, and the lack of accurate resource information leads to inefficient scheduling. In addition, and in contrast to the workloads of tightly coupled high-performance computing systems, a large part of the workloads submitted to these distributed systems consists of large sets (bags) of sequential tasks. Therefore, a realistic performance analysis of scheduling bags-of-tasks in large-scale distributed systems is important. Towards this end, we introduce in this paper a realistic workload model for bags-of-tasks, and we explore through trace-based simulations the design space of scheduling bags-of-tasks. Finally, we identify three new scheduling policies that use only inaccurate information when scheduling, and we compare them against known classes of proposed scheduling policies.
Alexandru Iosup, Omer Ozan Sonmez, Shanny Anoep, Dick H. J. Epema
HPDC1
2008 Efficient management of data center resources for massively multiplayer online games
abstract
Today's massively multiplayer online games (MMOGs) can include millions of concurrent players spread across the world. To keep these highly-interactive virtual environments online, a MMOG operator may need to provision tens of thousands of computing resources from various data centers. Faced with large resource demand variability, and with misfit resource renting policies, the current industry practice is to maintain for each game tens of self-owned data centers. In this work we investigate the dynamic resource provisioning from external data centers for MMOG operation. We introduce a novel MMOG workload model that represents the dynamics of both the player population and the player interactions. We evaluate several algorithms, including a novel neural network predictor, for predicting the resource demand. Using trace-based simulation, we evaluate the impact of the data center policies on the resource provisioning efficiency; we show that dynamic provisioning can be much more efficient than its static alternative.
Vlad Nae, Alexandru Iosup, Stefan Podlipnig, Radu Prodan, Dick H. J. Epema, Thomas Fahringer
SC2
2008 TRIBLER: a social-based peer-to-peer system
abstract
Abstract Most current peer‐to‐peer (P2P) file‐sharing systems treat their users as anonymous, unrelated entities, and completely disregard any social relationships between them. However, social phenomena such as friendship and the existence of communities of users with similar tastes or interests may well be exploited in such systems in order to increase their usability and performance. In this paper we present a novel social‐based P2P file‐sharing paradigm that exploits social phenomena by maintaining social networks and using these in content discovery, content recommendation, and downloading. Based on this paradigm's main concepts such as taste buddies and friends, we have designed and implemented the TRIBLER P2P file‐sharing system as a set of extensions to BitTorrent. We present and discuss the design of TRIBLER, and we show evidence that TRIBLER enables fast content discovery and recommendation at a low additional overhead, and a significant improvement in download performance. Copyright © 2007 John Wiley & Sons, Ltd.
Johan A. Pouwelse, Pawel Garbacki, Jun Wang 0012, Arno Bakker, Jie Yang 0015, Alexandru Iosup, Dick H. J. Epema, Marcel J. T. Reinders, Maarten van Steen, Henk J. Sips
Concurr. Comput. Pract. Exp.6
2008 The Grid Workloads Archive
Alexandru Iosup, Hui Li 0025, Mathieu Jan, Shanny Anoep, Catalin Dumitrescu, Lex Wolters, Dick H. J. Epema
Future Gener. Comput. Syst.1
2007 Build-and-Test Workloads for Grid Middleware: Problem, Analysis, and Applications
abstract
The Grid promise is starting to materialize today: large- scale multi-site infrastructures have grown to assist the work of scientists from all around the world. This tremendous growth can be sustained and continued only through a higher quality of the middleware, in terms of deployability and of correct functionality. A potential solution to this problem is the adoption of industry practices regarding middleware building and testing. However, it is unclear what good build-and-test environments for grid middleware should look like, and how to use them efficiently. In this work we address both these problems. First, we study the characteristics of the NMI build-and-test environment, which handles millions of testing tasks annually, for major Grid middleware such as Condor, Globus, VDT, and gLite. Through the analysis of a system-wide trace covering the past two years we find the main characteristics of the workload, as well as the performance of the system under load. Second, we propose mechanisms for more efficient test management and operation, and for resource provisioning and evaluation. Notably, we propose a generic test optimization technique that reduces the test time by 95%, while achieving 93% of the maximum accuracy, under real conditions.
Alexandru Iosup, Dick H. J. Epema
CCGRID1
2007 The Characteristics and Performance of Groups of Jobs in Grids
Alexandru Iosup, Mathieu Jan, Omer Ozan Sonmez, Dick H. J. Epema
Euro-Par1
2007 Inter-operating grids through delegated matchmaking
abstract
The grid vision of a single computing utility has yet to materíalize: while many grids with thousands of processors each exist, most work in isolation. An important obstacle for the effective and efficient inter-operation of grids is the problem of resource selection. In this paper we propose a solution to this problem that combines the hierarchical and decentralized approaches for interconnecting grids. In our solution, a hierarchy of grid sites is augmented with peer-to-peer connections between sites under the same administrative control. To operate this architecture, we employ the key concept of delegated matchmaking, which temporarily binds resources from remote sites to the local environment. With trace-based simulations we evaluate our solution under various infrastructural and load conditions, and we show that it outperforms other approaches to inter-operating grids. Specifically, we show that delegated matchmaking achieves up to 60% more goodput and completes 26% more jobs than its best alternative.
Alexandru Iosup, Dick H. J. Epema, Todd Tannenbaum, Matthew Farrellee, Miron Livny
SC1
2006 GRENCHMARK: A Framework for Analyzing, Testing, and Comparing Grids
abstract
Grid computing is becoming the natural way to aggregate and share large sets of heterogeneous resources. With the infrastructure becoming ready for the challenge, current grid development and acceptance hinge on proving that grids reliably support real applications, and on creating adequate benchmarks to quantify this support. However, grid applications are just beginning to emerge, and traditional benchmarks have yet to prove representative in grid environments. To address this chicken-and-egg problem, we propose a middle-way approach: create and run synthetic grid workloads comprising applications representative for today’s grids. For this purpose, we have designed and implemented GRENCHMARK, a framework for synthetic workload generation and submission. The framework greatly facilitates synthetic workload modeling, comes with over 35 synthetic and real applications, and is extensible and flexible. We show how the framework can be used for grid system analysis, functionality testing in grid environments, and for comparing different grid settings, and present the results obtained with GRENCHMARK in our multi-cluster grid, the DAS
Alexandru Iosup, Dick H. J. Epema
CCGRID1
2006 Correlating Topology and Path Characteristics of Overlay Networks and the Internet
Alexandru Iosup, Pawel Garbacki, Johan A. Pouwelse, Dick H. J. Epema
CCGRID1
2006 Provisioning and Scheduling Resources for World-Wide Data-Sharing Services
abstract
Grid computing is becoming the natural way to aggregate and share large and heterogeneous sets of resources. However, grid development and acceptance hinge on proving that grids reliably support large communities of users, and their real applications. In this paper we assess the ability of existing grid infrastructures to provision resources for a class of applications with numerous potential users, namely the class of world-wide data-sharing services. For this purpose, we first analyze the requirements of this class of applications, and match them against the existing spare capacity in three existing large-scale grid environments, namely OSG/Grid3, NorduGrid, and CERN LCG. Having shown that the existing capacity is insufficient, we devise and assess through trace-based simulation five domainspecific scheduling policies. Our findings give evidence that grid technology could be successfully leveraged for worldwide data-sharing services, without impacting the level of service for the currently existing load.
Alexandru Iosup, Pawel Garbacki, Dick H. J. Epema
e-Science1
2006 On Grid Performance Evaluation Using Synthetic Workloads
Alexandru Iosup, Dick H. J. Epema, Carsten Franke 0001, Alexander Papaspyrou, Lars Schley, Baiyi Song, Ramin Yahyapour
JSSPP1
2006 2Fast : Collaborative Downloads in P2P Networks
abstract
P2P systems that rely on the voluntary contribution of bandwidth by the individual peers may suffer from free riding. To address this problem, mechanisms enforcing fairness in bandwidth sharing have been designed, usually by limiting the download bandwidth to the available upload bandwidth. As in real environments the latter is much smaller than the former, these mechanisms severely affect the download performance of most peers. In this paper we propose a system called 2Fast, which solves this problem while preserving the fairness of bandwidth sharing. In 2Fast, we form groups of peers that collaborate in downloading a file on behalf of a single group member, which can thus use its full download bandwidth. A peer in our system can use its currently idle bandwidth to help other peers in their ongoing downloads, and get in return help during its own downloads. We assess the performance of 2Fast analytically and experimentally, the latter in both real and simulated environments. We find that in realistic bandwidth limit settings, 2Fast improves the download speed by up to a factor of 3.5 in comparison to state-of-the-art P2P download protocols
Pawel Garbacki, Alexandru Iosup, Dick H. J. Epema, Maarten van Steen
Peer-to-Peer Computing2