VLDB 2026 Research / reviewers in the wild / expert
Shadi Ibrahim
dblp:00/981
· DBLP profile ↗
56ranked-venue papers
7as first author
15since 2021 · last 2026
0000-0002-4306-5280ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 45 · 4 first-author · 11 since 2021Databases, data management, data science and information retrieval · 3 · 1 since 2021Artificial intelligence and machine learning · 2Computer networks · 2 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 2 · 1 since 2021Software engineering, systems software and programming languages · 1 · 1 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | ALTOCUMULUS: Enabling Efficient Erasure Coding in IPFS
Mohammad Rizk, Shadi Ibrahim, Thomas Lambert |
CCGrid | 2 |
| 2026 | THEIA: Scalable and Effective Data Prefetching for Geo-Distributed Erasure-Coded Storage
Marc Tranzer, Shadi Ibrahim |
ICDCS | 2 |
| 2026 | CFDGraph: Privacy-Preserving Graph Processing for Large-Scale Collaborative Fraud DetectionabstractInternational audience Qiulin Wu, Amelie Chi Zhou, Tristan Allard, Shadi Ibrahim, Yuhong Feng, Lichun Li, Amr El Abbadi |
ICDE | 4 |
| 2025 | EDDE: Container Deployment Framework Beyond the CloudabstractContainers, renowned for their lightweight nature and flexibility, have seen growing adoption for deploying edge services such as web applications. However, existing cloud-oriented container deployment frameworks fail to address the unique challenges of edge environments, including geographical distribution, device heterogeneity, and resource constraints. This oversight leads to suboptimal performance for latency-sensitive edge services like HPC/AI-powered autonomous driving and edge gaming, which demand rapid startup and immediate responsiveness. Hao Fan 0006, Shadi Ibrahim, Lin Gu 0002, Song Wu 0001 |
SC | 3 |
| 2025 | KubeSPT: Stateful Pod Teleportation for Service Resilience With Live MigrationabstractContainer orchestration systems, such as Kubernetes, streamline containerized application deployment. As more and more applications are being deployed in Kubernetes, there is an increasing need for rescheduling - relocating a running pod to different nodes - due to system upgrades, node failures, and load-balancing optimizations. Live migration, which transfers services from source nodes to target nodes with minimal downtime, is the ideal support for rescheduling. However, implementing live migration for pods that run stateful services is challenging, because Kubernetes manages pods as stateless. First, the current pod's network namespace initialization process causes a mismatch in the network state between the migrated pod and internal containers. Second, migrating the memory state results in extended downtime. Third, Kubernetes operations on pods do not consider preserving the state of the pods. Therefore, we propose KubeSPT to achieve live migration of stateful pods in rescheduling scenarios. Firstly, we synchronize the network state of pods and internal containers by controlling packet flow and implement fast service redirection. Secondly, we introduce a Hot Data and Lazy-Restore method for memory restoration to reduce migration downtime. Finally, we decouple pod migration operations from other Kubernetes operations to ensure compatibility with live migration. Experimental results show that KubeSPT reduces downtime by 86%-93% compared to current rescheduling methods. Hansheng Zhang, Song Wu 0001, Hao Fan 0006, Weibin Xue, Chen Yu 0003, Shadi Ibrahim, Hai Jin 0001 |
IEEE Trans. Serv. Comput. | 7 |
| 2024 | Toward Stream Processing Elasticity in Realistic Geo-Distributed EnvironmentsabstractStream data processing is a widely used technology for analysing IoT-generated data shortly after being produced, and delivering timely insights about them. Executing such analysis in geo-distributed platforms enables shorter delays between data production and processing and fewer disturbances due to potential instability of long-distance networks, while retaining the ability to scale the processing capacity up and down according to the demand. However, current stream processing systems were designed for environments made of homogeneous servers connected together using high-speed network links. We experimentally study the performance of Apache Flink coupled with the Gesscale auto-scaler in conditions which resemble those of geo-distributed platforms. We demonstrate that Flink’s backpressure mechanism should not be used as the only trigger for rescaling operations in heterogeneous network conditions. Raw performance, as well as performance predictability, also degrade quickly in the presence of stateful data processing operators and/or high network latency between the processing nodes. Khaled Arsalane, Guillaume Pierre, Shadi Ibrahim |
IC2E | 3 |
| 2024 | StreamBox: A Lightweight GPU SandBox for Serverless Inference Workflow
Hao Wu 0010, Junxiao Deng, Shadi Ibrahim, Song Wu 0001, Hao Fan 0006, Ziyue Cheng, Hai Jin 0001 |
USENIX ATC | 4 |
| 2024 | QoS-pro: A QoS-enhanced Transaction Processing Framework for Shared SSDsabstractSolid State Drives (SSDs) are widely used in data-intensive scenarios due to their high performance and decreasing cost. However, in shared environments, concurrent workloads can interfere with each other, leading to a violation of Quality of Service (QoS). While QoS mechanisms like fairness guarantees and latency constraints have been integrated into SSDs, existing transaction processing frameworks offer limited QoS guarantees and can significantly degrade overall performance in a shared environment. The reason is that the internal components of an SSD, originally designed to exploit parallelism, struggle to coordinate effectively when QoS mechanisms are applied to them. This article proposes a novel QoS -enhanced transaction pro cessing framework, called QoS-pro, which enhances QoS guarantees for concurrent workloads while maintaining high parallelism for SSDs. QoS-pro achieves this by redesigning transaction processing procedures to fully exploit the parallelism of shared SSDs and enhancing QoS-oriented transaction translation and scheduling with parallelism features in mind. In terms of fairness guarantees, QoS-pro outperforms state-of-the-art methods by achieving 96% fairness improvement and 64% maximum latency reduction. QoS-pro also shows almost no loss in throughput when compared with parallelism-oriented methods. Additionally, QoS-pro triggers the fewest Garbage Collection (GC) operations and minimally affects concurrently running workloads during GC operations. Hao Fan 0006, Yiliang Ye, Shadi Ibrahim, Xingru Li, Weibin Xue, Song Wu 0001, Chen Yu 0003, Xuanhua Shi, Hai Jin 0001 |
ACM Trans. Archit. Code Optim. | 3 |
| 2023 | QoS-Aware and Cost-Efficient Dynamic Resource Allocation for Serverless ML WorkflowsabstractMachine Learning (ML) workflows are increasingly deployed on serverless computing platforms to benefit from their elasticity and fine-grain pricing. Proper resource allocation is crucial to achieve fast and cost-efficient execution of serverless ML workflows (specially for hyperparameter tuning and model training). Unfortunately, existing resource allocation methods are static, treat functions equally, and rely on offline prediction, which limit their efficiency. In this paper, we introduce CE-scaling – a Cost-Efficient autoscaling framework for serverless ML work-flows. During the hyperparameter tuning, CE-scaling partitions resources across stages according to their exact usage to minimize resource waste. Moreover, it incorporates an online prediction method to dynamically adjust resources during model training. We implement and evaluate CE-scaling on AWS Lambda using various ML models. Evaluation results show that compared to state-of-the-art static resource allocation methods, CE-scaling can reduce the job completion time and the monetary cost by up to 63% and 41% for hyperparameter tuning, respectively; and by up to 58% and 38% for model training. Hao Wu 0010, Junxiao Deng, Hao Fan 0006, Shadi Ibrahim, Song Wu 0001, Hai Jin 0001 |
IPDPS | 4 |
| 2022 | Stragglers' Detection in Big Data Analytic Systems: The Impact of Heartbeat ArrivalabstractSpeculative execution can significantly improve the performance of Big Data applications by launching other copies of stragglers (slow tasks). Stragglers detection plays an important role in the effectiveness of speculative execution. The methods employed to detect stragglers use the information extracted from the last received heartbeats which may be outdated when triggering detection. This, in turn, can mislead Big Data analytic systems to make wrong detection with high inaccuracy. To shed the light on this issue, we carry out extensive simulations to identify how heartbeat arrival, task starting times, and detection methods impact the accuracy of stragglers detection in Big Data analytic systems. We reveal that the asynchrony in heartbeat arrivals not only lead to marking normal tasks as stragglers (false positives) but can also result in overlooking real stragglers (false negatives). Thomas Lambert, Shadi Ibrahim, Twinkle Jain, David Guyon |
CCGRID | 2 |
| 2022 | PGPregel: an end-to-end system for privacy-preserving graph processing in geo-distributed data centersabstractGraph processing is a popular computing model for big data analytics. Emerging big data applications are often maintained in multiple geographically distributed (geo-distributed) data centers (DCs) to provide low-latency services to global users. Graph processing in geo-distributed DCs suffers from costly inter-DC data communications. Furthermore, due to increasing privacy concerns, geo-distribution imposes diverse, strict, and often asymmetric privacy regulations that constrain geo-distributed graph processing. Existing graph processing systems fail to address these two challenges. In this paper, we design and implement PGPregel, which is an end-to-end system that provides privacy-preserving graph processing in geo-distributed DCs with low latency and high utility. To ensure privacy, PGPregel smartly integrates Differential Privacy into graph processing systems with the help of two core techniques, namely sampling and combiners, to reduce the amount of inter-DC data transfer while preserving good accuracy of graph processing results. We implement our design in Giraph and evaluate it in real cloud DCs. Results show that PGPregel can preserve the privacy of graph data with low overhead and good accuracy. Amelie Chi Zhou, Ruibo Qiu, Thomas Lambert, Tristan Allard, Shadi Ibrahim, Amr El Abbadi |
SoCC | 5 |
| 2022 | Container-aware I/O stack: bridging the gap between container storage drivers and solid state devicesabstractSolid State Devices (SSDs) have been widely adopted in containerized cloud platforms as they provide parallel and high-speed data accesses for critical data-intensive applications. Unfortunately, the I/O stack of the physical host overlooks the layered and independent nature of containers, thus I/O operations require expensive file redirect (between the storage driver, Overlay2/EXT4, and the virtual file system, VFS) and are scheduled sequentially. Moreover, containers suffer from significant I/O contention as resources at the native file system are shared between them. This paper presents a Container-aware I/O stack (CAST). CAST is made up of Layer-aware VFS (LaVFS) and Container-aware Native File System (CaFS). LaVFS locates files based on layer information and enables simultaneous Copy-on-Write (CoW) operations and thus avoids the overhead of searching and modifying files. CaFS, on the other hand, provides contention-free access by designing fine-grain resource allocation at the native file system. Experimental results using a NVMe SSD with micro-benchmarks and real-world applications show that CAST achieves 216%-219% (38%-98%, respectively) improvement over the original I/O stack. Song Wu 0001, Hao Fan 0006, Shadi Ibrahim, Hai Jin 0001 |
VEE | 5 |
| 2022 | Shadow: Exploiting the Power of Choice for Efficient Shuffling in MapReduceabstractHow to reduce the costly cross-rack data transferring is challenging in improving the performance of MapReduce platforms. Previous schemes mainly exploit the data locality in the Map phase to reduce the cross-rack communications. However, the Map locality based schemes may lead to highly skewed distribution of Map tasks across racks in the platform, resulting in serious load imbalance among different cross-rack links during Shuffling. Recent research results show that the slow Shuffling is the root cause of the MapReduce performance degradation. Very limited work has been done for speeding up the Shuffle phase. A notable scheme leverages the principle of the power of choice to balance the network loads on different cross-rack links during Shuffling for a specific type of sampling applications, where processing a random subset of the large-scale data collection is sufficient to derive the final result. The scheme launches a few additional tasks to offer more choices for task selection during Shuffling. However, such a scheme is designed for sampling applications and not applicable to general applications, where all the input data instead of a random subset is processed. In this work, we observe that with high Map locality, the network is mainly saturated in Shuffling but relatively free in the Map phase. A little sacrifice in Map locality may greatly accelerate Shuffling. Based on this, we propose a novel scheme called Shadow for Shuffle-constrained general applications, which strikes a trade-off between Map locality and Shuffling load balance. Specifically, Shadow iteratively chooses an original Map task from the most heavily loaded rack and creates a duplicated task for it on the most lightly loaded rack. During processing, Shadow makes a choice between an original task and its replica by efficiently pre-estimating the job execution time. We conduct extensive experiments to evaluate the Shadow design. Results show that Shadow greatly reduces the cross-rack skewness by 36.6 percent and the job execution time by 26 percent compared to existing schemes. Sijie Wu, Hanhua Chen, Hai Jin 0001, Shadi Ibrahim |
IEEE Trans. Big Data | 4 |
| 2022 | Taming System Dynamics on Resource Optimization for Data Processing Workflows: A Probabilistic ApproachabstractIn many data-intensive applications, workflow is often used as an important model for organizing data processing tasks and resource provisioning is an important and challenging problem for improving the performance of workflows. Recently, system variations in the cloud and large-scale clusters, such as those in I/O and network performances and failure events, have been observed to greatly affect the performance of workflows. Traditional resource provisioning methods, which overlook these variations, can lead to suboptimal resource provisioning results. In this article, we provide a general solution for workflow performance optimizations considering system variations. Specifically, we model system dynamics as time-dependent random variables and take their probability distributions as optimization input. Despite its effectiveness, this solution involves heavy computation overhead. Thus, we propose three pruning techniques to simplify workflow structure and reduce the probability evaluation overhead. We implement our techniques in a runtime library, which allows users to incorporate efficient probabilistic optimization into existing resource provisioning methods. Experiments show that probabilistic solutions can improve the performance by up to 65 percent compared to state-of-the-art static solutions, and our pruning techniques can greatly reduce the overhead of our probabilistic approach. Amelie Chi Zhou, Weilin Xue, Bingsheng He, Shadi Ibrahim, Reynold Cheng |
IEEE Trans. Parallel Distributed Syst. | 5 |
| 2021 | Gear: Enable Efficient Container Storage and Deployment with a New Image FormatabstractContainers have been widely used in various cloud platforms as they enable agile and elastic application deployment through their process-based virtualization and layered image system. However, different layers of a container image may contain substantial duplicate and unnecessary data, which slows down its deployment due to long image downloading time and increased burden on the image registry. To accelerate the deployment and reduce the size of the registry, we propose a new image format, named Gear image, that consists of two parts: a Gear index describing the structure of the image's file system and a set of files that are required when running an application. The Gear index is represented as a single-layer image compatible with the existing deployment framework. Containers can be launched by pulling a Gear index and on demand retrieving files pointed to by the index. Furthermore, the Gear image enables a file-level sharing mechanism, which helps remove duplicate data in the registry and avoid repeated downloading of identical files by a client. We implement a prototype of the container framework, named Gear, supporting the new image format. Evaluation shows that Gear saves 54 % storage capacity in the registry, speeds up container startup by up to${5\times}$, and reduces 84 % bandwidth demands. Hao Fan 0006, Shengwei Bian, Song Wu 0001, Song Jiang 0001, Shadi Ibrahim, Hai Jin 0001 |
ICDCS | 5 |
| 2020 | Rethinking Operators Placement of Stream Data Application in the EdgeabstractMaximum Sustainable Throughput (MST) refers to the amount of data that a Data Stream Processing (DSP) system can ingest while keeping stable performance. It has been acknowledged as an accurate metric to evaluate the performance of stream data processing. Yet, existing operators placements continue to focus on latency and throughput, not MST, as main performance objective when deploying stream data applications in the Edge. In this paper, we argue that MST should be used as an optimization objective when placing operators. This is specially important in the Edge, where network bandwidth and data streams are highly dynamic. We demonstrate that through the design and evaluation of a MST-driven operators placement (based on constraint programming) for stream data applications. Through simulations, we show how existing placement strategies that target overall communications reduction often fail to keep up with the rate of data streams. Importantly, the constraint programming-based operators placement is able to sustain up to 5x increased data ingestion compared to baseline strategies. Thomas Lambert, David Guyon, Shadi Ibrahim |
CIKM | 3 |
| 2020 | Cost-Aware Partitioning for Efficient Large Graph Processing in Geo-Distributed DatacentersabstractGraph processing is an emerging computation model for a wide range of applications and graph partitioning is important for optimizing the cost and performance of graph processing jobs. Recently, many graph applications store their data on geo-distributed datacenters (DCs) to provide services worldwide with low latency. This raises new challenges to existing graph partitioning methods, due to the multi-level heterogeneities in network bandwidth and communication prices in geo-distributed DCs. In this article, we propose an efficient graph partitioning method named Geo-Cut, which takes both the cost and performance objectives into consideration for large graph processing in geo-distributed DCs. Geo-Cut adopts two optimization stages. First, we propose a cost-aware streaming heuristic and utilize the one-pass streaming graph partitioning method to quickly assign edges to different DCs while minimizing inter-DC data communication cost. Second, we propose two partition refinement heuristics which identify the performance bottlenecks of geo-distributed graph processing and refine the partitioning result obtained in the first stage to reduce the inter-DC data transfer time while satisfying the budget constraint. Geo-Cut can be also applied to partition dynamic graphs thanks to its lightweight runtime overhead. We evaluate the effectiveness and efficiency of Geo-Cut using real-world graphs with both real geo-distributed DCs and simulations. Evaluation results show that Geo-Cut can reduce the inter-DC data transfer time by up to 79 percent (42 percent as the median) and reduce the monetary cost by up to 75 percent (26 percent as the median) compared to state-of-the-art graph partitioning methods with a low overhead. Amelie Chi Zhou, Bingkun Shen, Shadi Ibrahim, Bingsheng He |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2019 | On the Importance of Container Image Placement for Service Provisioning in the EdgeabstractEdge computing promises to extend Clouds by moving computation close to data sources to facilitate short-running and low-latency applications and services. Providing fast and predictable service provisioning time presets a new and mounting challenge, as the scale of Edge-servers grows and the heterogeneity of networks between them increases. This paper is driven by a simple question: can we place container images across Edge-servers in such a way that an image can be retrieved to any Edge-server fast and in a predictable time. To this end, we present KCBP and KCBP-WC, two container image placement algorithms which aim to reduce the maximum retrieval time of container images. KCBP and KCBP-WC are based on k-Center optimization. However, KCBP-WC tries to avoid placing large layers of a container image on the same Edge-server. Evaluations using trace-driven simulations show that KCBP and KCBP-WC can be applied to various network configurations and reduce the maximum retrieval time of container images by 1.1x to 4x compared to state-of-the-art placements (i.e., Best-Fit and Random). Jad Darrous, Thomas Lambert, Shadi Ibrahim |
ICCCN | 3 |
| 2019 | When FPGA-Accelerator Meets Stream Data Processing in the EdgeabstractToday, stream data applications represent the killer applications for Edge computing: placing computation close to the data source facilitates real-time analysis. Previous efforts have focused on introducing light-weight distributed stream processing (DSP) systems and dividing the computation between Edge servers and the clouds. Unfortunately, given the limited computation power of Edge servers, current efforts may fail in practice to achieve the desired latency of stream data applications. In this vision paper, we argue that by introducing FPGAs in Edge servers and integrating them into DSP systems, we might be able to realize stream data processing in Edge infrastructures. We demonstrate that through the design, implementation, and evaluation of F-Storm, an FPGA-accelerated and general-purpose distributed stream processing system on Edge servers. F-Storm integrates PCIe-based FPGAs into Edge-based stream processing systems and provides accelerators as a service for stream data applications. We evaluate F-Storm using different representative stream data applications. Our experiments show that compared to Storm, F-Storm reduces the latency by 36% and 75% for matrix multiplication and grep application. It also obtains 1.4x and 2.1x improvement for these two applications, respectively. We expect this work to accelerate progress in this domain. Song Wu 0001, Shadi Ibrahim, Hai Jin 0001, Jiang Xiao 0001, Haikun Liu |
ICDCS | 3 |
| 2019 | Incorporating Probabilistic Optimizations for Resource Provisioning of Data Processing WorkflowsabstractWorkflow is an important model for big data processing and resource provisioning is crucial to the performance of workflows. Recently, system variations in the cloud and large-scale clusters, such as those in I/O and network performances, have been observed to greatly affect the performance of workflows. Traditional resource provisioning methods, which overlook these variations, can lead to suboptimal resource provisioning results. In this paper, we provide a general solution for workflow performance optimizations considering system variations. Specifically, we model system variations as time-dependent random variables and take their probability distributions as optimization input. Despite its effectiveness, this solution involves heavy computation overhead. Thus, we propose three pruning techniques to simplify workflow structure and reduce the probability evaluation overhead. We implement our techniques in a runtime library, which allows users to incorporate efficient probabilistic optimization into existing resource provisioning methods. Experiments show that probabilistic solutions can improve the performance by 51% compared to state-of-the-art static solutions while guaranteeing budget constraint, and our pruning techniques can greatly reduce the overhead of probabilistic optimization. Amelie Chi Zhou, Bingsheng He, Shadi Ibrahim, Reynold Cheng |
ICPP | 4 |
| 2019 | NCQ-Aware I/O Scheduling for Conventional Solid State DrivesabstractWhile current fairness-driven I/O schedulers are successful in allocating equal time/resource share to concurrent workloads, they ignore the I/O request queueing or reordering in storage device layer, such as Native Command Queueing (NCQ). As a result, requests of different workloads cannot have an equal chance to enter NCQ (NCQ conflict) and fairness is violated. We address this issue by providing the first systematic empirical analysis on how NCQ affects I/O fairness and SSD utilization and accordingly proposing a NCQ-aware I/O scheduling scheme, NASS. The basic idea of NASS is to elaborately control the request dispatch of workloads to relieve NCQ conflict and improve NCQ utilization. NASS builds on two core components: an evaluation model to quantify important features of the workload, and a dispatch control algorithm to set the appropriate request dispatch of running workloads. We integrate NASS into four state-of-the-art I/O schedulers and evaluate its effectiveness using widely used benchmarks and real world applications. Results show that with NASS, I/O schedulers can achieve 11-23% better fairness and at the same time improve device utilization by 9-29%. Hao Fan 0006, Song Wu 0001, Shadi Ibrahim, Hai Jin 0001, Jiang Xiao 0001, Haibing Guan |
IPDPS | 3 |
| 2019 | Is it Time to Revisit Erasure Coding in Data-Intensive Clusters?abstractData-intensive clusters are heavily relying on distributed storage systems to accommodate the unprecedented growth of data. Hadoop distributed file system (HDFS) is the primary storage for data analytic frameworks such as Spark and Hadoop. Traditionally, HDFS operates under replication to ensure data availability and to allow locality-aware task execution of data-intensive applications. Recently, erasure coding (EC) is emerging as an alternative method to replication in storage systems due to the continuous reduction in its computation overhead. In this work, we conduct an extensive experimental study to understand the performance of data-intensive applications under replication and EC. We use representative benchmarks on the Grid'5000 testbed to evaluate how analytic workloads, data persistency, failures, the back-end storage devices, and the network configuration impact their performances. Our study sheds the light not only on the potential benefits of erasure coding in data-intensive clusters but also on the aspects that may help to realize it effectively. Jad Darrous, Shadi Ibrahim, Christian Pérez |
MASCOTS | 2 |
| 2018 | Nitro: Network-Aware Virtual Machine Image Management in Geo-Distributed CloudsabstractRecently, most large cloud providers, like Amazon and Microsoft, replicate their Virtual Machine Images (VMIs) on multiple geographically distributed data centers to offer fast service provisioning. Provisioning a service may require to transfer a VMI over the wide-area network (WAN) and therefore is dictated by the distribution of VMIs and the network bandwidth in-between sites. Nevertheless, existing methods to facilitate VMI management (i.e., retrieving VMIs) overlook network heterogeneity in geo-distributed clouds. In this paper, we design, implement and evaluate Nitro, a novel VMI management system that helps to minimize the transfer time of VMIs over a heterogeneous WAN. To achieve this goal, Nitro incorporates two complementary features. First, it makes use of deduplication to reduce the amount of data which will be transferred due to the high similarities within an image and in-between images. Second, Nitro is equipped with a network-aware data transfer strategy to effectively exploit links with high bandwidth when acquiring data and thus expedites the provisioning time. Experimental results show that our network-aware data transfer strategy offers the optimal solution when acquiring VMIs while introducing minimal overhead. Moreover, Nitro outperforms state-of-the-art VMI storage systems (e.g., OpenStack Swift) by up to 77%. Jad Darrous, Shadi Ibrahim, Amelie Chi Zhou, Christian Pérez |
CCGrid | 2 |
| 2018 | TurboStream: Towards Low-Latency Data Stream ProcessingabstractData Stream Processing (DSP) applications are often modelled as a directed acyclic graph: operators with data streams among them. Inter-operator communications can have a significant impact on the latency of DSP applications, accounting for 86% of the total latency. Despite their impact, there has been relatively little work on optimizing inter-operator communications, focusing on reducing inter-node traffic but not considering inter-process communication (IPC) inside a node, which often generates high latency due to the multiple memory-copy operations. This paper describes the design and implementation of TurboStream, a new DSP system designed specifically to address the high latency caused by inter-operator communications. To achieve this goal, we introduce (1) an improved IPC framework with OSRBuffer, a DSP-oriented buffer, to reduce memory-copy operations and waiting time of each single message when transmitting messages between the operators inside one node, and (2) a coarse-grained scheduler that consolidates operator instances and assigns them to nodes to diminish the inter-node IPC traffic. Using a prototype implementation, we show that our improved IPC framework reduces the end-to-end latency of intra-node IPC by 45.64% to 99.30%. Moreover, TurboStream reduces the latency of DSP by 83.23% compared to JStorm. Song Wu 0001, Shadi Ibrahim, Hai Jin 0001, Lin Gu 0002, Zhiyi Liu |
ICDCS | 3 |
| 2018 | Dual-Paradigm Stream ProcessingabstractExisting stream processing frameworks operate either under data stream paradigm processing data record by record to favor low latency, or under operation stream paradigm processing data in micro-batches to desire high throughput. For complex and mutable data processing requirements, this dilemma brings the selection and deployment of stream processing frameworks into an embarrassing situation. Moreover, current data stream or operation stream paradigms cannot handle data burst efficiently, which probably results in noticeable performance degradation. This paper introduces a dual-paradigm stream processing, called DO (Data and Operation) that can adapt to stream data volatility. It enables data to be processed in micro-batches (i.e., operation stream) when data burst occurs to achieve high throughput, while data is processed record by record (i.e., data stream) in the remaining time to sustain low latency. DO embraces a method to detect data bursts, identify the main operations affected by the data burst and switch paradigms accordingly. Our insight behind DO's design is that the trade-off between latency and throughput of stream processing frameworks can be dynamically achieved according to data communication among operations in a fine-grained manner (i.e., operation level) instead of framework level. We implement a prototype stream processing framework that adopts DO. Our experimental results show that our framework with DO can achieve 5x speedup over operation stream under low data stream sizes, and outperforms data stream on throughput by 2.1x to 3.2x under data burst. Song Wu 0001, Zhiyi Liu, Shadi Ibrahim, Lin Gu 0002, Hai Jin 0001 |
ICPP | 3 |
| 2018 | Energy-Efficient Speculative Execution using Advanced Reservation for Heterogeneous ClustersabstractMany Big Data processing applications nowadays run on large-scale multi-tenant clusters. Due to hardware heterogeneity and resource contentions, straggler problem has become the norm rather than the exception in such clusters. To handle the straggler problem, speculative execution has emerged as one of the most widely used straggler mitigation techniques. Although a number of speculative execution mechanisms have been proposed, as we have observed from real-world traces, the questions of "when" and "where" to launch speculative copies have not been fully discussed and hence cause inefficiencies on the performance and energy of Big Data applications. In this paper, we propose a performance model and an energy consumption model to reveal the performance and energy variations with different speculative execution solutions. We further propose a window-based dynamic resource reservation and a heterogeneity-aware copy allocation technique to answer the "when" and "where" questions for speculative executions. Evaluations using real-world traces show that our proposed technique can improve the performance of Big Data applications by up to 30% and reduce the overall energy consumption by up to 34%. Amelie Chi Zhou, Tien-Dat Phan, Shadi Ibrahim, Bingsheng He |
ICPP | 3 |
| 2018 | Improving the Effectiveness of Burst Buffers for Big Data Processing in HPC Systems with Eley
Orcun Yildiz, Amelie Chi Zhou, Shadi Ibrahim |
Future Gener. Comput. Syst. | 3 |
| 2017 | An Empirical Evaluation of How The Network Impacts The Performance and Energy Efficiency in RAMCloudabstractIn-memory storage systems emerged as a de-facto building block for today's large scale Web architectures and Big Data processing frameworks. Many research and engineering efforts have been dedicated to improve their performance and memory efficiency. More recently, such systems can leverage high-performance networks, e.g., Infiniband. To be able to leverage these systems, it is essential to understand the trade-offs induced by the use of high-performance networks. This paper aims to provide empirical evidence of the impact of client's location on the performance and energy consumption of in-memory storage systems. Through a study carried on RAMCloud, we focus on two settings: 1) clients are collocated within the same network as the storage servers (with Infiniband interconnects), 2) clients access the servers from a remote network, through TCP/IP. We compare and discuss aspects related to scalability and power consumption for these two scenarios which correspond to different deployment models for applications making use of in-memory cloud storage systems. Yacine Taleb, Shadi Ibrahim, Gabriel Antoniu, Toni Cortes |
CCGrid | 2 |
| 2017 | Eley: On the Effectiveness of Burst Buffers for Big Data Processing in HPC SystemsabstractBurst Buffer is an effective solution for reducing the data transfer time and the I/O interference in HPC systems. Extending Burst Buffers (BBs) to handle Big Data applications is challenging because BBs must account for the large data inputs of Big Data applications and the performance guarantees of HPC applications - which are considered as first-class citizens in HPC systems. Existing BBs focus on only intermediate data of Big Data applications and incur a high performance degradation of both Big Data and HPC applications. We present Eley, a burst buffer solution that helps to accelerate the performance of Big Data applications while guaranteeing the performance of HPC applications. In order to improve the performance of Big Data applications, Eley employs a prefetching technique that fetches the input data of these applications to be stored close to computing nodes thus reducing the latency of reading data inputs. Moreover, Eley is equipped with a full delay operator to guarantee the performance of HPC applications - as they are running independently on a HPC system. The experimental results show the effectiveness of Eley in obtaining shorter execution time of Big Data applications (shorter map phase) while guaranteeing the performance of HPC applications. Orcun Yildiz, Amelie Chi Zhou, Shadi Ibrahim |
CLUSTER | 3 |
| 2017 | Energy-Driven Straggler Mitigation in MapReduce
Tien-Dat Phan, Shadi Ibrahim, Amelie Chi Zhou, Guillaume Pallez, Gabriel Antoniu |
Euro-Par | 2 |
| 2017 | Characterizing Performance and Energy-Efficiency of the RAMCloud Storage SystemabstractMost large popular web applications, like Facebook and Twitter, have been relying on large amounts of in-memory storage to cache data and offer a low response time. As the main memory capacity of clusters and clouds increases, it becomes possible to keep most of the data in the main memory. This motivates the introduction of in-memory storage systems. While prior work has focused on how to exploit the low-latency of in-memory access at scale, there is very little visibility into the energy-efficiency of in-memory storage systems. Even though it is known that main memory is a fundamental energy bottleneck in computing systems (i.e., DRAM consumes up to 40% of a server's power). In this paper, by the means of experimental evaluation, we have studied the performance and energy-efficiency of RAMCloud - a well-known in-memory storage system. We reveal that although RAMCloud is scalable for read-only applications, it exhibits non-proportional power consumption. We also find that the current replication scheme implemented in RAMCloud limits the performance and results in high energy consumption. Surprisingly, we show that replication can also play a negative role in crash-recovery. Yacine Taleb, Shadi Ibrahim, Gabriel Antoniu, Toni Cortes |
ICDCS | 2 |
| 2017 | On Achieving Efficient Data Transfer for Graph Processing in Geo-Distributed DatacentersabstractGraph partitioning is important for optimizing the performance and communication cost of large graph processing jobs. Recently, many graph applications such as social networks store their data on geo-distributed datacenters (DCs) to provide services worldwide with low latency. This raises new challenges to existing graph partitioning methods, due to the costly Wide Area Network (WAN) usage and the multi-levels of network heterogeneities in geo-distributed DCs. In this paper, we propose a geo-aware graph partitioning method named G-Cut, which aims at minimizing the inter-DC data transfer time of graph processing jobs in geo-distributed DCs while satisfying the WAN usage budget. G-Cut adopts two novel optimization phases which address the two challenges in WAN usage and network heterogeneities separately. G-Cut can be also applied to partition dynamic graphs thanks to its light-weight runtime overhead. We evaluate the effectiveness and efficiency of G-Cut using realworld graphs with both real geo-distributed DCs and simulations. Evaluation results show that G-Cut can reduce the inter-DC data transfer time by up to 58% and reduce the WAN usage by up to 70% compared to state-of-the-art graph partitioning methods with a low runtime overhead. Amelie Chi Zhou, Shadi Ibrahim, Bingsheng He |
ICDCS | 2 |
| 2017 | Enabling fast failure recovery in shared Hadoop clusters: Towards failure-aware scheduling
Orcun Yildiz, Shadi Ibrahim, Gabriel Antoniu |
Future Gener. Comput. Syst. | 2 |
| 2016 | On the Root Causes of Cross-Application I/O Interference in HPC Storage SystemsabstractAs we move toward the exascale era, performance variability in HPC systems remains a challenge. I/O interference, a major cause of this variability, is becoming more important every day with the growing number of concurrent applications that share larger machines. Earlier research efforts on mitigating I/O interference focus on a single potential cause of interference (e.g., the network). Yet the root causes of I/O interference can be diverse. In this work, we conduct an extensive experimental campaign to explore the various root causes of I/O interference in HPC storage systems. We use microbenchmarks on the Grid'5000 testbed to evaluate how the applications' access pattern, the network components, the file system's configuration, and the backend storage devices influence I/O interference. Our studies reveal that in many situations interference is a result of bad flow control in the I/O path, rather than being caused by some single bottleneck in one of its components. We further show that interference-free behavior is not necessarily a sign of optimal performance. To the best of our knowledge, our work provides the first deep insight into the role of each of the potential root causes of interference and their interplay. Our findings can help developers and platform owners improve I/O performance and motivate further research addressing the problem across all components of the I/O stack. Orcun Yildiz, Matthieu Dorier, Shadi Ibrahim, Robert B. Ross, Gabriel Antoniu |
IPDPS | 3 |
| 2016 | iShare: Balancing I/O performance isolation and disk I/O efficiency in virtualized environmentsabstractSummary Performance isolation has long been a challenging problem for disk resource allocation in virtualized environments. While there have been many researches working on I/O performance isolation and disk utilization, none of them addresses the I/O performance isolation and disk utilization as a whole. To this end, we investigate the impact of current disk I/O performance isolation schemes on disk I/O utilization. Interestingly, our studies report that current isolation schemes bring unnecessary disk idle and reduce the overall disk I/O performance because of ignoring the disk states and characteristics of requests. Accordingly, we propose an adaptive proportional‐share I/O scheduling framework, namediShare, in virtualized environments.iSharenot only ensures I/O performance isolation through proportionally allocating time slices according to the weights of virtual machines but also preserves high disk efficiency by detecting disk states and adaptively adjusting the time slice size based on characteristics of requests. We implement a prototype ofiShareon the Xen platform. The experimental results show thatiShareensures I/O performance isolation while improving disk I/O efficiency, compared withBlkio(i.e., the default I/O performance isolation method in Xen),iShareincreases disk I/O bandwidth by 58% and slightly improves the I/O performance isolation for the sequential write applications. Copyright © 2015 John Wiley & Sons, Ltd. Song Wu 0001, Songqiao Tao, Hao Fan 0006, Hai Jin 0001, Shadi Ibrahim |
Concurr. Comput. Pract. Exp. | 6 |
| 2016 | On the energy footprint of I/O management in Exascale HPC systems
Matthieu Dorier, Orcun Yildiz, Shadi Ibrahim, Anne-Cécile Orgerie, Gabriel Antoniu |
Future Gener. Comput. Syst. | 3 |
| 2016 | Governing energy consumption in Hadoop through CPU frequency scaling: An analysis
Shadi Ibrahim, Tien-Dat Phan, Alexandra Carpen-Amarie, Houssem Chihoub, Diana Moise, Gabriel Antoniu |
Future Gener. Comput. Syst. | 1 |
| 2016 | Using Formal Grammars to Predict I/O Behaviors in HPC: The Omnisc'IO ApproachabstractThe increasing gap between the computation performance of post-petascale machines and the performance of their I/O subsystem has motivated many I/O optimizations including prefetching, caching, and scheduling. In order to further improve these techniques, modeling and predicting spatial and temporal I/O patterns of HPC applications as they run has become crucial. In this paper we present Omnisc'IO, an approach that builds a grammar-based model of the I/O behavior of HPC applications and uses it to predict when future I/O operations will occur, and where and how much data will be accessed. To infer grammars, Omnisc'IO is based on StarSequitur, a novel algorithm extending Nevill-Manning's Sequitur algorithm. Omnisc'IO is transparently integrated into the POSIX and MPI I/O stacks and does not require any modification in applications or higher-level I/O libraries. It works without any prior knowledge of the application and converges to accurate predictions of any N future I/O operations within a couple of iterations. Its implementation is efficient in both computation time and memory footprint. Matthieu Dorier, Shadi Ibrahim, Gabriel Antoniu, Robert B. Ross |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2015 | Chronos: Failure-aware scheduling in shared Hadoop clustersabstractHadoop emerged as the de facto state-of-the-art system for MapReduce-based data analytics. The reliability of Hadoop systems depends in part on how well they handle failures. Currently, Hadoop handles machine failures by re-executing all the tasks of the failed machines (i.e., executing recovery tasks). Unfortunately, this elegant solution is entirely entrusted to the core of Hadoop and hidden from Hadoop schedulers. The unawareness of failures therefore may prevent Hadoop schedulers from operating correctly towards meeting their objectives (e.g., fairness, job priority) and can significantly impact the performance of MapReduce applications. This paper presents Chronos, a failure-aware scheduling strategy that enables an early yet smart action for fast failure recovery while still operating within a specific scheduler objective. Upon failure detection, rather than waiting an uncertain amount of time to get resources for recovery tasks, Chronos leverages a lightweight preemption technique to carefully allocate these resources. In addition, Chronos considers data locality when scheduling recovery tasks to further improve the performance. We demonstrate the utility of Chronos by combining it with Fifo and Fair schedulers. The experimental results show that Chronos recovers to a correct scheduling behavior within a couple of seconds only and reduces the job completion times by up to 55% compared to state-of-the-art schedulers. Orcun Yildiz, Shadi Ibrahim, Tran Anh Phuong, Gabriel Antoniu |
IEEE BigData | 2 |
| 2015 | Exploring Energy-Consistency Trade-Offs in Cassandra Cloud Storage SystemabstractApache Cassandra is an open-source cloud storage system that offers multiple types of operation-level consistency including eventual consistency with multiple levels of guarantees and strong consistency. It is being used by many data-center applications (e.g., Facebook and App Scale). Most existing research efforts have been dedicated to exploring trade-offs such as: consistency vs. Performance, consistency vs. Latency and consistency vs. Monetary cost. In contrast, a little work is focusing on the consistency vs. Energy trade-off. As power bills have become a substantial part of the monetary cost for operating a data-center, this paper aims to provide a clearer understanding of the interplay between consistency and energy consumption. Accordingly, a series of experiments have been conducted to explore the implication of different factors on the energy consumption in Cassandra. Our experiments have revealed a noticeable variation in the energy consumption depending on the consistency level. Furthermore, for a given consistency level, the energy consumption of Cassandra varies with the access pattern and the load exhibited by the application. This further analysis indicates that the uneven distribution of the load amongst different nodes also impacts the energy consumption in Cassandra. Finally, we experimentally compare the impact of four storage configuration and data partitioning policies on the energy consumption in Cassandra: interestingly, we achieve 23% energy saving when assigning 50% of the nodes to the hot pool for the applications with moderate ratio of reads and writes, while applying eventual (quorum) consistency. This study points to opportunities for future research on consistency-energy trade-offs and offers useful insight into designing energy-efficient techniques for cloud storage systems. Houssem Chihoub, Shadi Ibrahim, Gabriel Antoniu, María S. Pérez 0001, Luc Bougé |
SBAC-PAD | 2 |
| 2015 | Spatial Locality Aware Disk Scheduling in Virtualized EnvironmentabstractExploiting spatial locality, a key technique for improving disk I/O utilization and performance, faces additional challenges in the virtualized cloud because of the transparency feature of virtualization. This paper contributes a novel disk I/O scheduling framework, named Pregather, to improve disk I/O efficiency through exposure and exploitation of the special spatial locality in the virtualized environment, thereby improving the performance of disk-intensive applications without harming the transparency feature of virtualization. The key idea behind Pregatheris to implement an intelligent model to predict the access regularity of spatial locality for each VM. Moreover, Pregather embraces an adaptive time slice allocation scheme to further reduce the resource contention and ensure fairness among VMs. We implement the Pregather disk scheduling framework and perform extensive experiments that involve multiple simultaneous applications of both synthetic benchmarks and MapReduce applications on Xen-based platforms. Our experiments demonstrate the accuracy of our prediction model and indicate that Pregather results in the high disk spatial locality and a significant improvement in disk throughput and application performance. Shadi Ibrahim, Song Wu 0001, Hai Jin 0001 |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2014 | CALCioM: Mitigating I/O Interference in HPC Systems through Cross-Application CoordinationabstractUnmatched computation and storage performance in new HPC systems have led to a plethora of I/O optimizations ranging from application-side collective I/O to network and disk-level request scheduling on the file system side. As we deal with ever larger machines, the interference produced by multiple applications accessing a shared parallel file system in a concurrent manner becomes a major problem. Interference often breaks single-application I/O optimizations, dramatically degrading application I/O performance and, as a result, lowering machine wide efficiency. This paper focuses on CALCioM, a framework that aims to mitigate I/O interference through the dynamic selection of appropriate scheduling policies. CALCioM allows several applications running on a supercomputer to communicate and coordinate their I/O strategy in order to avoid interfering with one another. In this work, we examine four I/O strategies that can be accommodated in this framework: serializing, interrupting, interfering and coordinating. Experiments on Argonne's BG/P Surveyor machine and on several clusters of the French Grid'5000 show how CALCioM can be used to efficiently and transparently improve the scheduling strategy between two otherwise interfering applications, given specified metrics of machine wide efficiency. Matthieu Dorier, Gabriel Antoniu, Robert B. Ross, Dries Kimpe, Shadi Ibrahim |
IPDPS | 5 |
| 2014 | Omnisc'IO: A Grammar-Based Approach to Spatial and Temporal I/O Patterns PredictionabstractThe increasing gap between the computation performance of post-petascale machines and the performance of their I/O subsystem has motivated many I/O optimizations including prefetching, caching, and scheduling techniques. In order to further improve these techniques, modeling and predicting spatial and temporal I/O patterns of HPC applications as they run has became crucial. In this paper we present Omnisc'IO, an approach that builds a grammar-based model of the I/O behavior of HPC applications and uses it to predict when future I/O operations will occur, and where and how much data will be accessed. Omnisc'IO is transparently integrated into the POSIX and MPI I/O stacks and does not require any modification in applications or higher level I/O libraries. It works without any prior knowledge of the application and converges to accurate predictions within a couple of iterations only. Its implementation is efficient in both computation time and memory footprint. Matthieu Dorier, Shadi Ibrahim, Gabriel Antoniu, Robert B. Ross |
SC | 2 |
| 2013 | Consistency in the Cloud: When Money Does Matter!abstractWith the emergence of cloud computing, many organizations have moved their data to the cloud in order to provide scalable, reliable and highly available services. To meet the ever-growing user needs, these services mainly rely on geographically-distributed data replication to guarantee good performance and high availability. However, with replication, consistency comes into question. Service providers in the cloud have the freedom to select the level of consistency according to the access patterns exhibited by the applications. Most optimizations efforts then concentrate on how to provide adequate trade-offs between consistency guarantees and performance. However, as the monetary cost completely relies on the service providers, in this paper we argue that monetary cost should be taken into consideration when evaluating or selecting a consistency level in the cloud. Accordingly, we define a new metric called consistency-cost efficiency. Based on this metric, we present a simple, yet efficient economical consistency model, called Bismar, that adaptively tunes the consistency level at runtime in order to reduce the monetary cost while simultaneously maintaining a low fraction of stale reads. Experimental evaluations with the Cassandra cloud storage on the Grid'5000 test bed show the validity of the metric and demonstrate the effectiveness of the proposed consistency model. Houssem Chihoub, Shadi Ibrahim, Gabriel Antoniu, María S. Pérez 0001 |
CCGRID | 2 |
| 2013 | Exploiting Spatial Locality to Improve Disk Efficiency in Virtualized EnvironmentsabstractVirtualization has become a prominent tool in data centers and is extensively leveraged in cloud environments: it enables multiple virtual machines (VMs) - with multiple operating systems and applications - to run within a physical server. However, virtualization introduces the challenging issue of preserving the high disk utilization (i.e., reducing the seek delay and rotation overhead) when allocating disk resources to VMs. Exploiting spatial locality, a key technique for improving disk utilization and performance, faces additional challenges in the virtualized cloud because of the transparency feature of virtualization (hyper visors do not have the information about the access patterns of applications running within each VM). To this end, this paper contributes a novel disk I/O scheduling framework, named Pregather, to improve disk I/O efficiency through exposure and exploitation of the special spatial locality in the virtualized environment (regional and sub-regional spatial locality corresponds to the virtual disk space and applications' access patterns, respectively), thereby improving the performance of disk-intensive applications without harming the transparency feature of virtualization (without a priori knowledge of the applications' access patterns). The key idea behind Pregather is to implement an intelligent model to predict the access regularity of sub-regional spatial locality for each VM. We implement the Pregather disk scheduling framework and perform extensive experiments that involve multiple simultaneous applications of both synthetic benchmarks and a MapReduce application on Xen-based platforms. Our experiments demonstrate the accuracy of our prediction model and indicate that Pregather results in the high disk spatial locality and a significant improvement in disk throughput and application performance. Shadi Ibrahim, Hai Jin 0001, Song Wu 0001, Songqiao Tao |
MASCOTS | 2 |
| 2013 | Flubber: Two-level disk scheduling in virtualized environment
Hai Jin 0001, Shadi Ibrahim, Wenzhi Cao, Song Wu 0001, Gabriel Antoniu |
Future Gener. Comput. Syst. | 3 |
| 2013 | Handling partitioning skew in MapReduce using LEEN
Shadi Ibrahim, Hai Jin 0001, Lu Lu 0006, Bingsheng He, Gabriel Antoniu, Song Wu 0001 |
Peer-to-Peer Netw. Appl. | 1 |
| 2013 | Petri net based Grid workflow verification and optimization
Haijun Cao, Hai Jin 0001, Song Wu 0001, Shadi Ibrahim |
J. Supercomput. | 4 |
| 2012 | Maestro: Replica-Aware Map Scheduling for MapReduceabstractMapReduce has emerged as a leading programming model for data-intensive computing. Many recent research efforts have focused on improving the performance of the distributed frameworks supporting this model. Many optimizations are network-oriented and most of them mainly address the data shuffling stage of MapReduce. Our studies with Hadoop demonstrate that, apart from the shuffling phase, another source of excessive network traffic is the high number of map task executions which process remote data. That leads to an excessive number of useless speculative executions of map tasks and to an unbalanced execution of map tasks across different machines. All these factors produce a noticeable performance degradation. We propose a novel scheduling algorithm for map tasks, named Maestro, to improve the overall performance of the MapReduce computation. Maestro schedules the map tasks in two waves: first, it fills the empty slots of each data node based on the number of hosted map tasks and on the replication scheme for their input data, second, runtime scheduling takes into account the probability of scheduling a map task on a given machine depending on the replicas of the task's input data. These two waves lead to a higher locality in the execution of map tasks and to a more balanced intermediate data distribution for the shuffling phase. In our experiments on a 100-node cluster, Maestro achieves around 95% local map executions, reduces speculative map tasks by 80% and results in an improvement of up to 34% in the execution time. Shadi Ibrahim, Hai Jin 0001, Lu Lu 0006, Bingsheng He, Gabriel Antoniu, Song Wu 0001 |
CCGRID | 1 |
| 2012 | Efficient Disk I/O Scheduling with QoS Guarantee for Xen-based Hosting PlatformsabstractIn this paper, we address the problem of allocating disk resources to guarantee specified latency and throughput targets of VMs while keeping efficient disk I/O. Accordingly, we present two-level scheduling framework, namely Flubber, in Xen-based hosting platform that decouples latency and throughput allocation. The high-level throughput control regulates the pending requests from the VMs, in order to meet the throughput requirements of different VMs and ensure isolation. Meanwhile, the low-level latency control, by the virtue of the batch and delay EDF mechanism, reorders all pending requests from VMs based on the their deadlines, and batches them to the disk device considering the locality of accesses across VMs. We have implemented Flubber with intensive evaluations on Xen-based host. The results show that Flubber can simultaneously meet the different service requirements of VMs while improving the efficiency of the physical disk. In contrast to CFQ, besides that Flubber achieves the desired QoS of each VM, Flubber speeds up the sequential and random read by 17% and 25% due to the efficient physical disk utilization. Hai Jin 0001, Shadi Ibrahim, Wenzhi Cao, Song Wu 0001 |
CCGRID | 3 |
| 2012 | Harmony: Towards Automated Self-Adaptive Consistency in Cloud StorageabstractIn just a few years cloud computing has become a very popular paradigm and a business success story, with storage being one of the key features. To achieve high data availability, cloud storage services rely on replication. In this context, one major challenge is data consistency. In contrast to traditional approaches that are mostly based on strong consistency, many cloud storage services opt for weaker consistency models in order to achieve better availability and performance. This comes at the cost of a high probability of stale data being read, as the replicas involved in the reads may not always have the most recent write. In this paper, we propose a novel approach, named Harmony, which adaptively tunes the consistency level at run-time according to the application requirements. The key idea behind Harmony is an intelligent estimation model of stale reads, allowing to elastically scale up or down the number of replicas involved in read operations to maintain a low (possibly zero) tolerable fraction of stale reads. As a result, Harmony can meet the desired consistency of the applications while achieving good performance. We have implemented Harmony and performed extensive evaluations with the Cassandra cloud storage on Grid'5000 test bed and on Amazon EC2. The results show that Harmony can achieve good performance without exceeding the tolerated number of stale reads. For instance, in contrast to the static eventual consistency used in Cassandra, Harmony reduces the stale data being read by almost 80% while adding only minimal latency. Meanwhile, it improves the throughput of the system by 45% while maintaining the desired consistency requirements of the applications when compared to the strong consistency model in Cassandra. Houssem Chihoub, Shadi Ibrahim, Gabriel Antoniu, María S. Pérez 0001 |
CLUSTER | 2 |
| 2011 | Adaptive Disk I/O Scheduling for MapReduce in Virtualized EnvironmentabstractVirtual machine (VM) interference has long been a challenging problem for performance predictability and system throughput for large-scale virtualized environments in the cloud. Such interferences are contributed by intertwined factors including the application's type, the number of con current VMs, and the VM scheduling algorithms used within the host. Since MapReduce has become an important data processing platform in the cloud, we investigate the impact of disk schedulers in Hadoop. Interestingly, our experimental results report a noticeable variation of the Hadoop performance between different applications when applying different disk pairs' schedulers in both the hypervisor and the virtual machines. Furthermore, a typical Hadoop application consists of different interleaving stages, each requiring different I/O workloads and patterns. As a result, the disk pairs' schedulers are not only sub-optimal for different MapReduce applications, but also sub-optimal for different sub-phases of the whole job. Accordingly, this paper presents a novel approach for adaptively tuning the disk pairs' schedulers in both the hypervisor and the virtual machines during the execution of a single MapReduce job. Our results show that MapReduce performance can be significantly improved; specifically, adaptive tuning of disk pairs' schedulers achieves a 25% performance improvement on a sort benchmark with Hadoop. Shadi Ibrahim, Hai Jin 0001, Lu Lu 0006, Bingsheng He, Song Wu 0001 |
ICPP | 1 |
| 2010 | LEEN: Locality/Fairness-Aware Key Partitioning for MapReduce in the CloudabstractThis paper investigates the problem of Partitioning Skew in MapReduce-based system. Our studies with Hadoop, a widely used MapReduce implementation, demonstrate that the presence of partitioning skew causes a huge amount of data transfer during the shuffle phase and leads to significant unfairness on the reduce input among different data nodes. As a result, the applications experience performance degradation due to the long data transfer during the shuffle phase along with the computation skew, particularly in reduce phase. We develop a novel algorithm named LEEN for locality-aware and fairness-aware key partitioning in MapReduce. LEEN embraces an asynchronous map and reduce scheme. All buffered intermediate keys are partitioned according to their frequencies and the fairness of the expected data distribution after the shuffle phase. We have integrated LEEN into Hadoop-0.18.0. Our experiments demonstrate that LEEN can efficiently achieve higher locality and reduce the amount of shuffled data. More importantly, LEEN guarantees fair distribution of the reduce inputs. As a result, LEEN achieves a performance improvement of up to 40% on different workloads. Shadi Ibrahim, Hai Jin 0001, Lu Lu 0006, Song Wu 0001, Bingsheng He |
CloudCom | 1 |
| 2010 | MR-scope: a real-time tracing tool for MapReduceabstractMapReduce programming model is emerging as an efficient tool for data-intensive applications. Hadoop, an open-source implementation of MapReduce, has been widely adopted and experienced by both academia and enterprise. Recently, lots of efforts have been done on improving the performance of MapReduce system and on analyzing the MapReduce process based on the log files generated during the Hadoop execution. Visualizing log files seems to be a very useful tool to understand the behavior of the Hadoop process. In this paper, we present MR-Scope, a real-time MapReduce tracing tool. MR-Scope provides a real-time insight of the MapReduce process, including the ongoing progress of every task hosted in Task Tracker. In addition, it displays the health of the Hadoop cluster data nodes, the distribution of the file system's blocks and their replicas and the content of the different block splits of the file system. We implement MR-Scope in native Hadoop 0.1. Experimental results demonstrat that MR-Scope's overhead is less than 4% when running wordcount benchmark. Dachuan Huang, Xuanhua Shi, Shadi Ibrahim, Lu Lu 0006, Hongzhang Liu, Song Wu 0001, Hai Jin 0001 |
HPDC | 3 |
| 2009 | Evaluating MapReduce on Virtual Machines: The Hadoop Case
Shadi Ibrahim, Hai Jin 0001, Lu Lu 0006, Song Wu 0001, Xuanhua Shi |
CloudCom | 1 |
| 2009 | CLOUDLET: towards mapreduce implementation on virtual machinesabstractThe existing MapReduce framework in virtualized environment suffers from poor performance, due to the heavy overhead of I/O virtualization, and management difficulty for storage and computation. To address the problems, we propose Cloudlet, a novel MapReduce framework on virtual machines. The aim of Cloudlet design is to overcome the overhead of VM while benefiting of the other features of VM (i.e. management and reliability issues). Shadi Ibrahim, Hai Jin 0001, Bin Cheng 0001, Haijun Cao, Song Wu 0001 |
HPDC | 1 |