VLDB 2026 Research / reviewers in the wild / expert
Bogdan Ghit
dblp:69/8814 · also Bogdan Ionut Ghit
· DBLP profile ↗
15ranked-venue papers
8as first author
2since 2021 · last 2022
0000-0002-2530-8736ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 10 · 7 first-author · 2 since 2021Software engineering, systems software and programming languages · 3 · 1 first-authorDatabases, data management, data science and information retrieval · 2Artificial intelligence and machine learning · 1Computer networks · 1 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 1
Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.
| Computer architecture, parallel and distributed computing, and storage systems
2 papers |
Cloud and datacenter computing · 42% Distributed systems · 41% Parallel and multicore computing · 12% | |
| Databases, data mining, and information retrieval
1 paper |
Information retrieval · 62% Query processing and optimization · 38% |
Topics — the 10 heaviest of 10, each with the papers that count most for it
| Topic | Weight | Papers | Last | Evidence papers |
|---|---|---|---|---|
Cloud and datacenter computing › cluster resource management and scheduling
cluster resource management |
0.5 | 2 | 2017 | Better Safe than Sorry: Grappling with Failures of In-Memory Data Analytics Frameworks · HPDC 2017 Balanced resource allocations across multiple dynamic MapReduce clusters · SIGMETRICS 2014 |
Information retrieval
indexing |
0.4 | 1 | 2019 | [Demo] Low-latency Spark Queries on Updatable Data · SIGMOD Conference 2019 |
Distributed systems › fault tolerance
checkpointing |
0.3 | 1 | 2017 | Better Safe than Sorry: Grappling with Failures of In-Memory Data Analytics Frameworks · HPDC 2017 |
Distributed systems
fault tolerance |
0.3 | 1 | 2017 | Better Safe than Sorry: Grappling with Failures of In-Memory Data Analytics Frameworks · HPDC 2017 |
Cloud and datacenter computing › resource allocation
dynamic resource allocation |
0.2 | 1 | 2014 | Balanced resource allocations across multiple dynamic MapReduce clusters · SIGMETRICS 2014 |
Parallel and multicore computing › data-parallel programming
mapreduce |
0.2 | 1 | 2014 | Balanced resource allocations across multiple dynamic MapReduce clusters · SIGMETRICS 2014 |
Query processing and optimization
interactive query processing |
0.1 | 1 | 2019 | [Demo] Low-latency Spark Queries on Updatable Data · SIGMOD Conference 2019 |
Query processing and optimization › query execution
low-latency query processing |
0.1 | 1 | 2019 | [Demo] Low-latency Spark Queries on Updatable Data · SIGMOD Conference 2019 |
Distributed systems › fault tolerance › checkpointing
checkpointing strategies |
0.1 | 1 | 2017 | Better Safe than Sorry: Grappling with Failures of In-Memory Data Analytics Frameworks · HPDC 2017 |
Storage systems
storage reliability |
0.1 | 1 | 2017 | Better Safe than Sorry: Grappling with Failures of In-Memory Data Analytics Frameworks · HPDC 2017 |
Methods — techniques the papers use, named apart from their topics
multi-version concurrency control · 0.4indexing · 0.4caching · 0.4weighting policies · 0.2monitoring · 0.2
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2022 | In-Memory Indexed Caching for Distributed Data ProcessingabstractPowerful abstractions such as dataframes are only as efficient as their underlying runtime system. The de-facto distributed data processing framework, Apache Spark, is poorly suited for the modern cloud-based data-science workloads due to its outdated assumptions: static datasets analyzed using coarse-grained transformations. In this paper, we introduce the Indexed DataFrame, an in-memory cache that supports a dataframe abstraction which incorporates indexing capabilities to support fast lookup and join operations. Moreover, it supports appends with multi-version concurrency control. We implement the Indexed DataFrame as a lightweight, standalone library which can be integrated with minimum effort in existing Spark programs. We analyze the performance of the Indexed DataFrame in cluster and cloud deployments with real-world datasets and benchmarks using both Apache Spark and Databricks Runtime. In our evaluation, we show that the Indexed DataFrame significantly speeds-up query execution when compared to a non-indexed dataframe, incurring modest memory overhead. Alexandru Uta, Bogdan Ghit, Ankur Dave, Jan S. Rellermeyer, Peter Boncz |
IPDPS | 2 |
| 2021 | Capri: Achieving Predictable Performance in Cloud Spot MarketsabstractLarge cloud providers offer spot instances at attractive prices to improve resource utilization, resulting in a spot market where users bid for resources and providers alter prices dynamically. As prices surpass bid values, resources may be relinquished from users with low bids. Achieving predictable performance on spot markets is challenging for data analytics workloads because they are very sensitive to preemptions due to the excessive cost of recomputations.We introduce capri, a scheduling system for running cloud data analytics in spot markets in which users may experience periods of degraded performance. capri dynamically predicts the functional relationship between bid and performance, thus helping with managing expectations and bid advice. We propose a new spot market abstraction called the bribe scheduler which delivers differentiated service levels based on bids. capri uses a prediction mechanism built on a queueing approximation of the bribe scheduler. capri dynamically estimates parameters to adapt the queueing model and provide accurate performance predictions in the face of time-varying workloads.We collect measurements using capri running two realistic workloads, imdb and tpcds, and demonstrate the accuracy of our approximation and parameter estimation methodology. We show that capri achieves a median prediction error below 3% in bursty workloads. We find that capri‘s service level prediction is pessimistic as users are likely to experience better performance than they should receive for their bids. Bogdan Ghit, Asser N. Tantawi |
MASCOTS | 1 |
| 2019 | [Demo] Low-latency Spark Queries on Updatable DataabstractAs data science gets deployed more and more into operational applications, it becomes important for data science frameworks to be able to perform computations in interactive, sub-second time. Indexing and caching are two key techniques that can make interactive query processing on large datasets possible. In this demo, we show the design, implementation and performance of a new indexing abstraction in Apache Spark, called the Indexed DataFrame. This is a cached DataFrame that incorporates an index to support fast lookup and join operations, and supports updates with multi-version concurrency. We demonstrate the Indexed Dataframe on a social network dataset using microbenchmarks and real-world graph processing queries, in datasets that are continuously growing. Alexandru Uta, Bogdan Ghit, Ankur Dave, Peter Boncz |
SIGMOD Conference | 2 |
| 2017 | Better Safe than Sorry: Grappling with Failures of In-Memory Data Analytics FrameworksabstractProviding fault-tolerance is of major importance for data analytics frameworks such as Hadoop and Spark, which are typically deployed in large clusters that are known to experience high failures rates. Unexpected events such as compute node failures are in particular an important challenge for in-memory data analytics frameworks, as the widely adopted approach to deal with them is to recompute work already done. Recomputing lost work, however, requires allocation of extra resource to re-execute tasks, thus increasing the job runtimes. To address this problem, we design a checkpointing system called Panda that is tailored to the intrinsic characteristics of data analytics frameworks. In particular, Panda employs fine-grained checkpointing at the level of task outputs and dynamically identifies tasks that are worthwhile to be checkpointed rather than be recomputed. As has been abundantly shown, tasks of data analytics jobs may have very variable runtimes and output sizes. These properties form the basis of three checkpointing policies which we incorporate into Panda. Bogdan Ghit, Dick H. J. Epema |
HPDC | 1 |
| 2017 | An Experimental Performance Evaluation of Autoscaling Policies for Complex WorkflowsabstractSimplifying 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 |
ICPE | 5 |
| 2016 | Tyrex: Size-Based Resource Allocation in MapReduce FrameworksabstractMany large-scale data analytics infrastructures are employed for a wide variety of jobs, ranging from short interactive queries to large data analysis jobs that may take hours or even days to complete. As a consequence, data-processing frameworks like MapReduce may have workloads consisting of jobs with heavy-tailed processing requirements. With such workloads, short jobs may experience slowdowns that are an order of magnitude larger than large jobs do, while the users may expect slowdowns that are more in proportion with the job sizes. To address this problem of large job slowdown variability in MapReduce frameworks, we design a scheduling system called TYREX that is inspired by the well-known TAGS task assignment policy in distributed-server systems. In particular, TYREX partitions the resources of a MapReduce framework, allowing any job running in any partition to read data stored on any machine, imposes runtime limits in the partitions, and successively executes parts of jobs in a work-conserving way in these partitions until they can run to completion. We develop a statistical model for dynamically setting the runtime limits that achieves near optimal job slowdown performance, and we empirically evaluate TYREX on a cluster system with workloads consisting of both synthetic and real-world benchmarks. We find that TYREX cuts in half the job slowdown variability while preserving the median job slowdown when compared to state-of-the-art MapReduce schedulers such as FIFO and FAIR. Furthermore, TYREX reduces the job slowdown at the 95th percentile by more than 50% when compared to FIFO and by 20-40% when compared to FAIR. Bogdan Ghit, Dick H. J. Epema |
CCGrid | 1 |
| 2016 | Which Cloud Auto-Scaler Should I Use for my Application?: Benchmarking Auto-Scaling AlgorithmsabstractRapid 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 |
ICPE | 3 |
| 2015 | Scheduling Workloads of Workflows with Unknown Task RuntimesabstractWorkflows are important computational tools in many branches of science, and because of the dependencies among their tasks and their widely different characteristics, scheduling them is a difficult problem. Most research on scheduling workflows has focused on the offline problem of minimizing the make span of single workflows with known task runtimes. The problem of scheduling multiple workflows has been addressed either in an offline fashion, or still with the assumption of known task runtimes. In this paper, we study the problem of scheduling workloads consisting of an arrival stream of workflows without task runtime estimates. The resource requirements of a workflow can significantly fluctuate during its execution. Thus, we present four scheduling policies for workloads of workflows with as their main feature the extent to which they reserve processors to workflows to deal with these fluctuations. We perform simulations with realistic synthetic workloads and we show that any form of processor reservation only decreases the overall system performance and that a greedy backfilling-like policy performs best. Alexey Ilyushkin, Bogdan Ghit, Dick H. J. Epema |
CCGRID | 2 |
| 2015 | Reducing Job Slowdown Variability for Data-Intensive WorkloadsabstractA well-known problem when executing data-intensive workloads with such frameworks as MapReduce is that small jobs with processing requirements counted in the minutes may suffer from the presence of huge jobs requiring hours or days of compute time, leading to a job slowdown distribution that is very variable and that is uneven across jobs of different sizes. Previous solutions to this problem for sequential or rigid jobs in single-server and distributed-server systems include priority-based FeedBack Queueing (FBQ), and Task Assignment by Guessing Sizes (TAGS), which kills and restarts from scratch on another server jobs that exceed the local time limit. In this paper, we derive four scheduling policies that are rightful descendants of existing size-based scheduling disciplines (among which FBQ and TAGS) with appropriate adaptations to data-intensive frameworks. The two main mechanisms employed by these policies are partitioning the resources of the datacenter, and isolating jobs with different size ranges. We evaluate these policies by means of realistic simulations of representative MapReduce workloads from Facebook and show that under the best of these policies, the vast majority of short jobs in MapReduce workloads experience close to ideal job slowdowns even under high system loads (in the range of 0.7-0.9) while the slowdown of the very large jobs is not prohibitive. We validate our simulations by means of experiments on a real multicluster system, and we find that the job slowdown performance results obtained with both match remarkably well. Bogdan Ghit, Dick H. J. Epema |
MASCOTS | 1 |
| 2014 | V for Vicissitude: The Challenge of Scaling Complex Big Data WorkflowsabstractIn 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 |
CCGRID | 1 |
| 2014 | KOALA-C: A task allocator for integrated multicluster and multicloud environmentsabstractCompanies, 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 |
CLUSTER | 2 |
| 2014 | Balanced resource allocations across multiple dynamic MapReduce clustersabstractRunning 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 |
SIGMETRICS | 1 |
| 2013 | The BTWorld use case for big data analytics: Description, MapReduce logical workflow, and empirical evaluationabstractThe 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 BigData | 2 |
| 2013 | Towards an Optimized Big Data Processing SystemabstractScalable 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 |
CCGRID | 1 |
| 2012 | Demonstrating BooSTER: The broadcast stream transmission epidemic repairabstractWireless broadcasting systems, such as Digital Video Broadcasting (DVB), are subject to signal degradation, having an effect on end users' reception quality. Reception quality can be improved by increasing signal strength, but this comes at a significantly increased energy use and still without guaranteeing error-free reception. BOOSTER is a system for cooperative repair of DVB streams. It operates based on a fully decentralized epidemic algorithm that boosts reception quality by cooperatively repairing lossy packet streams among the community of DVB viewers. In this demo, we demonstrate BOOSTER in a special configurable implementation that allows on-the-fly tweaking of the error rate, network topology, and some protocol parameters, allowing the audience to get a feeling on how they affect end-user experience and load distribution. Bogdan Ghit, Spyros Voulgaris, Aaron Harwood |
P2P | 1 |