EDBT 2026 Demo / reviewers in the wild / expert
Ningfang Mi
dblp:03/4718
· DBLP profile ↗
89ranked-venue papers
9as first author
17since 2021 · last 2026
0009-0005-0934-1287ORCID · corroborated
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 29 · 6 first-author · 9 since 2021Computer networks · 28 · 1 first-author · 1 since 2021Applied, interdisciplinary, general and emerging computing · 14 · 3 since 2021Software engineering, systems software and programming languages · 6 · 2 first-authorDatabases, data management, data science and information retrieval · 5Security and privacy · 4 · 2 first-authorGraphics, computer vision, multimedia, augmented reality and games · 4 · 1 since 2021Artificial intelligence and machine learning · 1Theory of computation · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | AFL-PRF: Adaptive Federated Learning for Low-Quality Data: Enhancing Performance, Robustness, and FairnessabstractFederated learning (FL) enables collaborative model training across distributed clients while preserving privacy, yet its decentralized nature makes it vulnerable to poisoned updates and performance degradation under highly skewed data. Prior studies typically treat accuracy, robustness, and fairness separately, leaving open the challenge of a unified solution. We propose AFL-PRF, an adaptive federated learning framework that simultaneously enhances accuracy, robustness, and fairness in adversarial and heterogeneous environments. AFL-PRF integrates three key techniques. First, an exponential adaptive weighting mechanism dynamically scales client updates, suppressing poisoned or unreliable contributions while retaining meaningful signals from benign but low-quality clients. Second, a client prioritization strategy guided by the novel Weight Update Divergence (WUD) score promotes reliable updates and their benign neighbors, preventing malicious gradients from dominating aggregation. Third, sensitivity profiling identifies fully connected (FC) layers as highly vulnerable due to large weight variance, motivating a selective clipping strategy that filters extreme updates in these layers while preserving normal learning dynamics. Extensive experiments on benchmark datasets demonstrate that AFL-PRF consistently outperforms state-of-the-art baselines, achieving over 30% improvement in robustness and 20% enhancement in fairness, while maintaining superior predictive accuracy. By unifying adaptive weighting, client prioritization, and targeted clipping, AFL-PRF establishes a new benchmark for federated learning under poisoned and highly non-IID conditions. Pinrui Yu, Longtian Ye, Geng Yuan, Ningfang Mi, Xue Lin 0001 |
WACV | 5 |
| 2025 | PECO: Probabilistic Evaluation-Based Client Selection for Federated Learning with Overlapping ClientsabstractFederated Learning (FL) enables privacy-preserving distributed machine learning by training models across clients without raw data exchange. However, non-IID data distributions-particularly when clients possess overlapping samples alongside unique data-pose significant challenges for client selection strategies. Traditional random sampling approaches inadequately balance the trade-off between redundant updates from overlapping data and diverse contributions from clientspecific samples, resulting in suboptimal model performance. We propose PECO, a dynamic client selection framework that evaluates each client's contribution to model performance using a small validation set and adaptively assigns selection probabilities based on these contributions. Through experiments on CIFAR10 and Fashion-MNIST, we demonstrate that PECO consistently outperforms random selection, achieving superior model accuracy and convergence stability across diverse non-IIID scenarios with overlapping data.11This work is partially supported by the National Science Foundation Award OCA-2417715 Shiyue Hou, Bo Sheng, Ningfang Mi |
HPCC | 5 |
| 2024 | Performance Analysis of Data Processing in Distributed File Systems with Near Data ProcessingabstractIn the era of big data, the escalating volume and velocity of data generation pose significant challenges in data processing. Traditional systems like Spark [1] and Hadoop [2] manage the increasing amount and velocity of data by improving data placement and processing speeds. However, they face inherent limitations due to the essential data movement required for processing. In this paper, we explore the Skyhook framework, a novel extension of the Ceph distributed system, which significantly reduces the need for data movement. We present an extensive case study using the Skyhook framework, applying it with the TPC-H and K-means clustering algorithms. More specifically, we leverage the TPC-H benchmark to distinguish between CPU-intensive and I/O-intensive tasks. We explore the integration of K-means clustering into SQL, coupled with a near-data processing system to offload the computational burden of the K-means clustering algorithm to storage nodes. We conduct a comprehensive performance evaluation of distributed data processing applications across three processing approaches: traditional layout (baseline), optimized layout, and near-data processing. Additionally, we introduce the use of the FIO tool to simulate real-world system workloads, enabling the measurement of performance metrics such as average latency and CPU utilization. Our research is a significant advance in understanding how to optimize data processing systems to meet the demands of the modern data landscape. Shiyue Hou, Nathan R. Tallent, Ningfang Mi |
ISNCC | 4 |
| 2024 | Uncovering the Impact of Bursty Workloads on System Performance in Serverless ComputingabstractServerless computing has emerged as a rapidly evolving paradigm within cloud services. Understanding the diverse arrival patterns of bursty serverless workload and their impact on performance and cost within serverless platforms is paramount. To the best of our knowledge, our study represents the first comprehensive analysis of observed bursty serverless workload characteristics, uncovering both anticipated and unexpected findings that illuminate the intricate interactions among workload characteristics, serverless platforms, and the associated performance-cost trade-offs. Through this analysis, we aim to furnish data-driven insights to serverless cloud providers, system administrators, and application developers, offering guidance on navigating the performance-cost trade-offs and potential pitfalls with various bursty workload arrival patterns in the realm of serverless computing. Shiyue Hou, Ningfang Mi |
ISNCC | 4 |
| 2024 | Adaptive Homogeneity-Based Client Selection Policy for Federated LearningabstractFederated learning (FL) is a distributed paradigm that enables multiple clients or edge devices to collaboratively train a model without sharing their local data. The FL system has to tackle a significant challenge due to the non-IID nature of client data. Traditional methods of client selection usually suffer from the problem of variable test accuracies as well as slow convergence because of their incapability in effectively handle data heterogeneity across clients. In this paper, we propose an Adaptive Homogeneity-Based Client Selection Policy (ASTraFL) to address this challenge. ASTraFL dynamically selects clients whose data distributions optimally complement the current state of the global model, focusing on increased homogeneity of the selected client data in each training round. Our experiments demonstrate that ASTraFL can accelerate convergence speed and ensure the learning process's robustness. Pinrui Yu, Geng Yuan, Xue Lin 0001, Ningfang Mi |
ISNCC | 5 |
| 2024 | A Data-Loader Tunable Knob to Shorten GPU Idleness for Distributed Deep LearningabstractDeep Neural Networks (DNNs) have been applied as an effective machine learning algorithm to tackle problems in different domains. However, the endeavor to train sophisticated DNN models can stretch from days into weeks, presenting substantial obstacles in the realm of research focused on large-scale DNN architectures. Distributed Deep Learning (DDL) contributes to accelerating DNN training by distributing training workloads across multiple computation accelerators, for example, graphics processing units (GPUs). Despite the considerable amount of research directed toward enhancing DDL training, the influence of data loading on GPU utilization and overall training efficacy remains relatively overlooked. It is non-trivial to optimize data-loading in DDL applications that need intensive central processing unit (CPU) and input/output (I/O) resources to process enormous training data. When multiple DDL applications are deployed on a system (e.g., Cloud and High-Performance Computing (HPC) system), the lack of a practical and efficient technique for data-loader allocation incurs GPU idleness and degrades the training throughput. Therefore, our work first focuses on investigating the impact of data-loading on the global training throughput. We then propose a throughput prediction model to predict the maximum throughput for an individual DDL training application. By leveraging the predicted results, A-Dloader is designed to dynamically allocate CPU and I/O resources to concurrently running DDL applications and use the data-loader allocation as a knob to reduce GPU idle intervals and thus improve the overall training throughput. We implement and evaluate A-Dloader in a DDL framework for a series of DDL applications arriving and completing across the runtime. Our experimental results show that A-Dloader can achieve a 28.9% throughput improvement and a 10% makespan improvement compared with allocating resources evenly across applications. Danlin Jia, Geng Yuan, Xue Lin 0001, Ningfang Mi |
ACM Trans. Archit. Code Optim. | 5 |
| 2024 | Learning-Based Dynamic Memory Allocation Schemes for Apache Spark Data ProcessingabstractApache Spark is an in-memory analytic framework that has been adopted in the industry and research fields. Two memory managers, Static and Unified, are available in Spark to allocate memory for caching Resilient Distributed Datasets (RDDs) and executing tasks. However, we find that the static memory manager (SMM) lacks flexibility, while the unified memory manager (UMM) puts heavy pressure on the garbage collection of the JVM on which Spark resides. To address these issues, we design a learning-based bidirectional usage-bounded memory allocation scheme to support dynamic memory allocation with the consideration of both memory demands and latency introduced by garbage collection. We first develop an auto-tuning memory manager (ATuMm) that adopts an intuitive feedback-based learning solution. However, ATuMm is a slow learner that can only alter the states of Java Virtual Memory (JVM) Heap in a limited range. That is, ATuMm decides to increase or decrease the boundary between the execution and storage memory pools by a fixed portion of JVM Heap size. To overcome this shortcoming, we further develop a new reinforcement learning-based memory manager (Q-ATuMm) that uses a Q-learning intelligent agent to dynamically learn and tune the partition of JVM Heap. We implement our new memory managers in Spark 2.4.0 and evaluate them by conducting experiments in a real Spark cluster. Our experimental results show that our memory manager can reduce the total garbage collection time and thus further improve Spark applications’ performance (i.e., reduced latency) compared to the existing Spark memory management solutions. By integrating our machine learning-driven memory manager into Spark, we can further obtain around 1.3x times reduction in the latency. Danlin Jia, Natalia Valencia, Janki Bhimani, Bo Sheng, Ningfang Mi |
IEEE Trans. Cloud Comput. | 6 |
| 2023 | MoKE: Modular Key-value Emulator for Realistic Studies on Emerging Storage DevicesabstractKey-value stores are widely used as building blocks in today's IT infrastructure for managing and storing large amounts of data. Storage technologies are undergoing continuous innovations to accelerate KV workloads. However, designing high-performance KV or object storage devices is challenging and still needs more research to address the performance bottlenecks of the existing designs. There is a void for an inexpensive and extendable research platform that enables in-depth exploration of the index management components within the KV devices. To fill this void, we design Modular Key-value Emulator (MoKE). MoKE is a software emulator for fostering future full-stack software/hardware KV and object storage device research. MoKE is cheap (software-based emulator), usable with SNIA KV API (supports popular host-device interfaces), extendable (supports internal KV device research), and adaptable (QEMU-based). Manoj Pravakar Saha, Danlin Jia, Janki Bhimani, Ningfang Mi |
CLOUD | 4 |
| 2023 | SRC: Mitigate I/O Throughput Degradation in Network Congestion Control of Disaggregated Storage SystemsabstractThe industry has adopted disaggregated storage systems to provide high-quality services for hyper-scale architectures. This infrastructure enables organizations to access storage resources that can be independently managed, configured, and scaled. It is supported by the recent advances of all-flash arrays and NVMe-over-Fabric protocol, enabling remote access to NVMe devices over different network fabrics. A surge of research has been proposed to mitigate network congestion in traditional remote direct memory access protocol (RDMA). However, NVMe-oF raises new challenges in congestion control for disaggregated storage systems.In this work, we investigate the performance degradation of the read throughput on storage nodes caused by traditional network congestion control mechanisms. We design a storage-side rate control (SRC) to relieve network congestion while avoiding performance degradation on storage nodes. First, we design an I/O throughput control mechanism in the NVMe driver layer to enable throughput control on storage nodes. Second, we construct a throughput prediction model to learn a mapping function between workload characteristics and I/O throughput. Third, we deploy SRC on storage nodes to cooperate with traditional network congestion control on an NVMe-over-RDMA architecture. Finally, we evaluate SRC with varying workloads, SSD configurations, and network topologies. The experimental results show that SRC achieves significant performance improvement. Danlin Jia, Xuebin Yao, Mahsa Bayati, Pradeep Subedi, Bo Sheng, Ningfang Mi |
IPDPS | 10 |
| 2022 | A Data-Loader Tunable Knob to Shorten GPU Idleness for Distributed Deep LearningabstractDeep Neural Network (DNN) has been applied as an effective machine learning algorithm to tackle problems in different domains. However, training a sophisticated DNN model takes days to weeks and becomes a challenge in constructing research on large-scale DNN models. Distributed Deep Learning (DDL) contributes to accelerating DNN training by distributing training workloads across multiple computation accelerators (e.g., GPUs). Although a surge of research works has been devoted to optimizing DDL training, the impact of data-loading on GPU usage and training performance has been relatively under-explored. It is non-trivial to optimize data-loading in DDL applications that need intensive CPU and I/O resources to process enormous training data. When multiple DDL applications are deployed on a system (e.g., Cloud and HPC), the lack of a practical and efficient technique for data-loader allocation incurs GPU idleness and degrades the training throughput. Therefore, our work first focuses on investigating the impact of data-loading on the global training throughput. We then propose a throughput prediction model to predict the maximum throughput for an individual DDL training application. By leveraging the predicted results, A-Dloader is designed to dynamically allocate CPU and I/O resources to concurrently running DDL applications and use the data-loader allocation as a knob to reduce GPU idle intervals and thus improve the overall training throughput. We implement and evaluate A-Dloader in a DDL framework for a series of DDL applications arriving and completing across the runtime. Our experimental results show that A-Dloader can achieve a 23.5% throughput improvement and a 10% makespan improvement, compared to allocating resources evenly across applications. Danlin Jia, Geng Yuan, Xue Lin 0001, Ningfang Mi |
CLOUD | 4 |
| 2022 | I/O Workload Management for All-Flash Datacenter Storage Systems Based on Total Cost of OwnershipabstractRecently, the capital expenditure of flash-based Solid State Driver (SSDs) keeps declining and the storage capacity of SSDs keeps increasing. As a result, all-flash storage systems have started to become more economically viable for large shared storage installations in datacenters, where metrics like Total Cost of Ownership (TCO) are of paramount importance. On the other hand, flash devices suffer from write amplification, which, if unaccounted, can substantially increase the TCO of a storage system. In this paper, we first develop a TCO model for datacenter all-flash storage systems, and then plug a Write Amplification model (WAF) of NVMe SSDs we build based on empirical data into this TCO model. Our new WAF model accounts for workload characteristics like write rate and percentage of sequential writes. Furthermore, using both the TCO and WAF models as the optimization criterion, we design new flash resource management schemes (minTCO) to guide datacenter managers to make workload allocation decisions under the consideration of TCO for SSDs. Based on that, we also developminTCO-RAIDto support RAID SSDs andminTCO-Offlineto optimize the offline workload-disk deployment problem during the initialization phase. Experimental results show thatminTCOcan reduce the TCO and keep relatively high throughput and space utilization of the entire datacenter storage resources. Zhengyu Yang 0001, Manu Awasthi, Mrinmoy Ghosh, Janki Bhimani, Ningfang Mi |
IEEE Trans. Big Data | 5 |
| 2022 | Auto-Tuning Parameters for Emerging Multi-Stream Flash-Based Storage Drives Through New I/O Pattern GenerationsabstractIn the era of big data processing, more and more data centers in cloud storage are now replacing traditional HDDs with enterprise SSDs. Both developers and users of these SSDs require thorough benchmarking to evaluate and configure the variable parameters of emerging technologies.[2]and[3]are the recent development of the SSD industry, which assists in placing data on SSDs in a smart way to improve application performance and SSD endurance. The challenging part to use multi-stream SSDs is to assign stream IDs to incoming writes, such that each stream consists of data with a similar lifetime. The benefit of the stream management algorithms varies over different workloads. Thus, first, we propose a new framework, calledPatternI/Ogenerator (PatIO), to capture the enterprise storage behavior that is prevailing across various user workloads, virtualization setup, file systems, and volume managers for the database server applications on flash-based storage. Second, usingPatIO, we study what type of applications may be benefited by which stream assignment algorithm. Third, we design the framework to automatically tune the variable parameters of different stream identification algorithms of the multi-stream SSDs. Our evaluation shows 20 to 110 percent of the reward function increase, measuring the cumulative impact on application performance and SSD endurance. Janki Bhimani, Adnan Maruf, Ningfang Mi, Rajinikanth Pandurangan, Vijay Balakrishnan |
IEEE Trans. Computers | 3 |
| 2022 | Automatic Stream Identification to Improve Flash Endurance in Data CentersabstractThe demand for high performance I/O in Storage-as-a-Service (SaaS) is increasing day by day. To address this demand, NAND Flash-based Solid-state Drives (SSDs) are commonly used in data centers as cache- or top-tiers in the storage rack ascribe to their superior performance compared to traditional hard disk drives (HDDs). Meanwhile, with the capital expenditure of SSDs declining and the storage capacity of SSDs increasing, all-flash data centers are evolving to serve cloud services better than SSD-HDD hybrid data centers. During this transition, the biggest challenge is how to reduce the Write Amplification Factor (WAF) as well as to improve the endurance of SSD since this device has a limited program/erase cycles. A specified case is that storing data with different lifetimes (i.e., I/O streams with similar temporal fetching patterns such as reaccess frequency) in one single SSD can cause high WAF, reduce the endurance, and downgrade the performance of SSDs. Motivated by this, multi-stream SSDs have been developed to enable data with a different lifetime to be stored in different SSD regions. The logic behind this is to reduce the internal movement of data—when garbage collection is triggered, there are high chances of having data blocks with either all the pages being invalid or valid. However, the limitation of this technology is that the system needs to manually assign the same streamID to data with a similar lifetime. Unfortunately, when data arrives, it is not known how important this data is and how long this data will stay unmodified. Moreover, according to our observation, with different definitions of a lifetime (i.e., different calculation formulas based on selected features previously exhibited by data, such as sequentiality, and frequency), streamID identification may have varying impacts on the final WAF of multi-stream SSDs. Thus, in this article, we first develop a portable and adaptable framework to study the impacts of different workload features and their combinations on write amplification. We then propose a feature-based stream identification approach, which automatically co-relates the measurable workload attributes (such as I/O size, I/O rate, and so on.) with high-level workload features (such as frequency, sequentiality, and so on.) and determines a right combination of workload features for assigning streamIDs . Finally, we develop an adaptable stream assignment technique to assign streamID for changing workloads dynamically. Our evaluation results show that our automation approach of stream detection and separation can effectively reduce the WAF by using appropriate features for stream assignment with minimal implementation overhead. Janki Bhimani, Zhengyu Yang 0001, Jingpei Yang, Adnan Maruf, Ningfang Mi, Rajinikanth Pandurangan, Changho Choi, Vijay Balakrishnan |
ACM Trans. Storage | 5 |
| 2021 | SNIS: Storage-Network Iterative Simulation for Disaggregated Storage SystemsabstractIn recent years, designs and optimizations on disaggregated storage systems supported by cutting-edge storage and network techniques emerge dramatically. However, conducting experiments in a disaggregated architecture is often expensive. A comprehensive modeling system of disaggregated storage systems is indispensable for researchers to construct fast and reliable experiments. Modeling and evaluating the performance of disaggregated storage systems is a challenge for the following reasons. First, the performance of a disaggregated storage system depends on network protocols and storage solutions jointly. Second, the available trace datasets for generating the workload may not suffice the need for the simulation, as they are collected without considering the integration of network delay and storage processing time. This work proposes a storage-network iterative simulation (SNIS) for disaggregated storage systems by considering the issues above. Our simulation methodology integrates storage and network simulations to model the end-to-end performance of disaggregated storage systems and conducts multiple rounds of simulations to update arrival times of read/write requests. The evaluation results show that SNIS can converge to a relatively stable state after a certain number of iterations. Danlin Jia, Tengpeng Li, Mahsa Bayati, Ron Lee, Bo Sheng, Ningfang Mi |
IPCCC | 8 |
| 2021 | Fine-grained control of concurrency within KV-SSDsabstractThe development of KV-SSDs allows simplifying the I/O stack compared to the traditional block-based SSDs. We propose a novel Key-Value-based Storage infrastructure for Parallel Computing(KV-SiPC)-a framework for multi-thread OpenMP applications to use NVMe-based KV-SSDs. We design a new capability to execute workloads with multiple parallel data threads along with traditional parallel compute threads, that allow us to improve the overall throughput of applications, utilizing the maximum possible storage bandwidth. We implement our KV-SiPC infrastructure in a real system by extending various processing layers (e.g., program, OS, and device layers) and evaluate the performance of KV-SiPC by using block-based NVMe SSDs in the traditional I/O stack as a baseline for comparisons. The experimental results show that KV-SiPC can better utilize the available device bandwidth and significantly increases application I/O throughput. Janki Bhimani, Jingpei Yang, Ningfang Mi, Changho Choi, Manoj Pravakar Saha, Adnan Maruf |
SYSTOR | 3 |
| 2021 | New YARN Non-Exclusive Resource Management Scheme through Opportunistic Idle Resource AssignmentabstractEfficiently managing resources and improving throughput in a large-scale cluster has become a crucial problem with the explosion of data processing applications in recent years. Hadoop YARN and Mesos, as two universal resource management platforms, have been widely adopted in the commodity cluster for co-deploying multiple data processing frameworks, such as Hadoop MapReduce and Apache Spark. However, in the existing resource management, a certain amount of resources are exclusively allocated to a running task and can only be re-assigned after that task is completed. This exclusive mode unfortunately leads to a potential problem that may under-utilize the cluster resources and degrade system performance. To address this issue, we propose a novel opportunistic and efficient resource allocation scheme, named OpERA, which breaks the barriers among the encapsulated resource containers by leveraging the knowledge of actual runtime resource utilizations to re-assign opportunistic available resources to the pending tasks. OpERA avoids incurring severe performance interference to active tasks by further using two approaches to efficiently balances the starvations of reserved tasks and normal queued tasks. We implement and evaluate OpERA in Hadoop YARN v2.5. Our experimental results show that OpERA significantly reduces the average job execution time and increases the resource (CPU and memory) utilizations. Zhengyu Yang 0001, Han Gao 0013, Ningfang Mi, Bo Sheng |
IEEE Trans. Cloud Comput. | 5 |
| 2021 | New Scheduling Algorithms for Improving Performance and Resource Utilization in Hadoop YARN ClustersabstractThe MapReduce framework has become the defacto scheme for scalable semi-structured and un-structured data processing in recent years. The Hadoop ecosystem has evolved into its second generation, Hadoop YARN, which adopts fine-grained resource management schemes for job scheduling. Nowadays, fairness and efficiency are two main concerns in YARN resource management because resources in YARN are shared and contended by multiple applications. However, the current scheduling in YARN does not yield the optimal resource arrangement, unnecessarily causing idle resources and inefficient scheduling. It omits the dependency between tasks which is extremely crucial for the efficiency of resource utilization as well as heterogeneous job features in real application environments. We thus propose a new YARN scheduler which can effectively reduce the makespan (i.e., the total execution time) of a batch of MapReduce jobs in Hadoop YARN clusters by leveraging the information of requested resources, resource capacities and dependency between tasks. For accommodating heterogeneity in MapReduce jobs, we also extend our scheduler by further considering the job iteration information in the scheduling decisions. We implemented the new scheduling algorithm as a pluggable scheduler in YARN and evaluated it with a set of classic MapReduce benchmarks. The experimental results demonstrate that our YARN scheduler effectively reduces the makespans and improves resource utilizations. Han Gao 0013, Bo Sheng, Ningfang Mi |
IEEE Trans. Cloud Comput. | 5 |
| 2020 | RITA: Efficient Memory Allocation Scheme for Containerized Parallel Systems to Improve Data Processing LatencyabstractThe emerging resource-sharing container-based virtualization is prevalent in IT, as it is a much lighter deployment in the cloud environment compared to VM-based virtualization. Distributed data-processing workloads executing in parallel take advantages of resource sharing, fast delivery, and excellent portability of containerization, but also suffer from resource competition and performance interference. Especially for memory virtualization, data-processing frameworks allocate physical memory (i.e., RAM) and swap to applications specified by users, without considering cache-characteristics and parallelism of applications, which induces performance degradation and significantly protracted latency which is worse given over-provisioning. We design an efficient memory allocation scheme (RITA) for containerized parallel systems to improve data processing latency. RITA monitors memory usage and cache characteristics of applications, and dynamically re-allocates memory resources. We implement RITA in a real-world system, which can easily migrate to other container-based virtualization environments. Our experimental results show that RITA provides remarkable latency improvement for memory intensive distributed data-processing workloads. Danlin Jia, Mahsa Bayati, Ron Lee, Ningfang Mi |
CLOUD | 4 |
| 2020 | Deploying Network Key-Value SSDs to Disaggregate Resources in Big Data Processing FrameworksabstractThe exponential data generation embraces unstructured object storage systems as an effective solution to improve performance. Key-Value (KV) SSD object storage devices are unveiled to mitigate the shortcomings of traditional Key-Value stores on block devices, including device low-bandwidth utilization and KV-store resource-draining operations on the host CPU and block devices. Samsung KV-SSDs are built on top of NVMe over Fabric hardware, which supports storage remote access protocols (i.e., RDMA). Network Key-Value (NKV) is a software eco-system developed by Samsung that enables data distribution and storage disaggregation of KV-SSDs. Most widely used big data processing platforms, such as Hadoop, Presto, deploy Hadoop Distributed File System (HDFS) to take advantage of rapid data access by co-locating storage and compute nodes. The co-allocation of compute and storage node limits the scalability and utilization resources and thus increases the total cost of ownership. In this paper, we present a new storage disaggregation model for big data processing platforms. Our new system layout leverages resource disaggregation by separating compute infrastructure from storage infrastructure and utilizes the benefits of new evolving storage technology, i.e., KV-SSD, for large-scale data access and processing. The goal of this work is to facilitate independent scaling of storage and compute resources, and shift the data retrieval load from the hosts to storage nodes. We evaluate our designed architecture using TPC-DS benchmark. Our results show that the CPU load on compute nodes is non-negligibly released with sustaining the same performance compared to the conventional Hadoop with HDFS. Mahsa Bayati, Harsh Roogi, Ron Lee, Ningfang Mi |
IPCCC | 4 |
| 2020 | Performance and Consistency Analysis for Distributed Deep Learning ApplicationsabstractAccelerating the training of Deep Neural Network (DNN) models is very important for successfully using deep learning techniques in fields like computer vision and speech recognition. Distributed frameworks help to speed up the training process for large DNN models and datasets. Plenty of works have been done to improve model accuracy and training efficiency, based on mathematical analysis of computations in the Con-volutional Neural Networks (CNN). However, to run distributed deep learning applications in the real world, users and developers need to consider the impacts of system resource distribution. In this work, we deploy a real distributed deep learning cluster with multiple virtual machines. We conduct an in-depth analysis to understand the impacts of system configurations, distribution typologies, and application parameters, on the latency and correctness of the distributed deep learning applications. We analyze the performance diversity under different model consistency and data parallelism by profiling run-time system utilization and tracking application activities. Based on our observations and analysis, we develop design guidelines for accelerating distributed deep-learning training on virtualized environments. Danlin Jia, Manoj Pravakar Saha, Janki Bhimani, Ningfang Mi |
IPCCC | 4 |
| 2019 | What does Vibration do to Your SSD?abstractVibration generated in modern computing environments such as autonomous vehicles, edge computing infrastructure, and data center systems is an increasing concern. In this paper, we systematically measure, quantify and characterize the impact of vibration on the performance of SSD devices. Our experiments and analysis uncover that exposure to both short-term and long-term vibration, even within the vendor-specified limits, can significantly affect SSD I/O performance and reliability. Janki Bhimani, Tirthak Patel, Ningfang Mi, Devesh Tiwari |
DAC | 3 |
| 2019 | Emulate Processing of Assorted Database Server Applications on Flash-Based Storage in Datacenter InfrastructuresabstractIn the era of big data processing, more and more datacenters in cloud storages are now replacing traditional HDDs with enterprise SSDs. Both developers and users of these SSDs require thorough benchmarking to evaluate their performance impacts. I/O performance with synthetic workload or classic benchmark varies drastically from real I/O activities in the datacenter. Thus, we propose a new framework, called Pattern I/O generator (PatIO), to collectively capture the enterprise storage behavior that is prevailing across assorted user workloads and system configurations for different database server applications on flash-based storage. PatIO is designed to emulate the processing of real-world I/O activities easily with less time and resource requirements. Our methodology comprises three main steps: (1) dissect the overall I/O activities of various real workloads and identify the prevailing attributes in distinct visual I/O patterns; (2) construct a pattern warehouse as the collection of unique I/O patterns that are generated through various combinations of multiple I/O jobs; and (3) finally integrate different combinations of these synthetically generated I/O patterns to reproduce the comprehensive characteristics of various real workloads and system setup for the database server applications. To provide an easy-to-use experience, we develop a graphical user interface (GUI). We evaluate our framework by comparing I/O characteristics and I/O performance of generated workloads with those of real-world workloads for multiple database applications such as MySQL, Cassandra, and ForestDB. Janki Bhimani, Rajinikanth Pandurangan, Ningfang Mi, Vijay Balakrishnan |
IPCCC | 3 |
| 2019 | ATuMm: Auto-tuning Memory Manager in Apache SparkabstractApache Spark is an in-memory analytic framework that has been adopted in the industry and research fields. Two memory managers, Static and Unified, are available in Spark to allocate memory for caching Resilient Distributed Datasets (RDDs) and executing tasks. However, we found that the static memory manager (SMM) lacks flexibility, while the unified memory manager (UMM) puts heavy pressure on the garbage collection of JVM on which Spark resides. To address these issues, we design an auto-tuning memory manager (ATuMm) to support dynamic memory allocation with the consideration of both memory demands and latency introduced by garbage collection. We implement our new memory manager in Spark 2.2.0 and evaluate it by conducting experiments in a real Spark cluster. Our experimental results show that our auto-tuning memory manager can reduce the total garbage collection time and thus further improve the performance (i.e., reduced latency) of Spark applications, compared to the existing Spark memory management solutions. Danlin Jia, Janki Bhimani, Son Nam Nguyen, Bo Sheng, Ningfang Mi |
IPCCC | 5 |
| 2019 | Abstract cost models for distributed data-intensive computations
Rundong Li 0001, Ningfang Mi, Mirek Riedewald, Yizhou Sun |
Distributed Parallel Databases | 2 |
| 2018 | BloomStream: Data Temperature Identification for Flash Based Memory Storage Using Bloom FiltersabstractData temperature identification is an importance issue of many fields like data caching and storage tiering in modern flash-based storage systems. With the technological advancement of memory and storage, data temperature identification is no longer just a classification of hot and cold, but instead becomes a "multistreaming" data categorization problem to classify data into multiple categories according to their temperature. Therefore, we propose a novel data temperature identification scheme that adopts bloom filters to efficiently capture both frequency and recency of data blocks and accurately identify the exact data temperature for each data block. Moreover, in bloom filter data structure we replace the original OR operation with the XOR masking operation such that our scheme can delete or reset bits in bloom filters and thus avoid high false positives due to saturation. We further utilize twin bloom filters to alternatively keep unmasked clean copies of data and thus ensure low false negative rate. Our extensive evaluation results show that our new scheme can accurately identify the exact data temperature with low false identification rates across different synthetic and real I/O workloads. More importantly, our scheme consumes less memory space compared to other existing data temperature identification schemes. Janki Bhimani, Ningfang Mi, Bo Sheng |
IEEE CLOUD | 2 |
| 2018 | FIOS: Feature Based I/O Stream Identification for Improving Endurance of Multi-Stream SSDsabstractThe demand for high speed 'Storage-as-a-Service' (SaaS) is increasing day-by-day. SSDs are commonly used in higher tiers of storage rack in data centers. Also, all flash data centers are evolving to better serve cloud services. Although SSDs guaranty better performance when compared to HDDs, but SSDs endurance is still a matter of concern. Storing data with different lifetime in an SSD can cause high write amplification and reduce the endurance and performance of SSDs. Recently, multi-stream SSDs have been developed to enable data with different lifetime to be stored in different SSD regions and thus reduce write amplification. To efficiently use this new multi-streaming technology, it is important to choose appropriate workload features to assign the same streamID to data with similar lifetime. However, we found that streamID identification using different features may have varying impacts on the final write amplification of multi-stream SSDs. Therefore, in this paper we develop a portable and adoptable framework to study the impacts of different workload features and their combinations on write amplification. We also introduce a new feature, named "coherency", to capture the friendship among write operations with respect to their update time. Finally, we propose a feature-based stream identification approach, which co-relates the measurable workload attributes (such as I/O size, I/O rate, etc.) with high level workload features (such as frequency, sequentiality etc.) and determines a good combination of workload features for assigning streamIDs. Our evaluation results show that our proposed approach can always reduce the Write Amplification Factor (WAF) by using appropriate features for stream assignment. Janki Bhimani, Ningfang Mi, Zhengyu Yang 0001, Jingpei Yang, Rajinikanth Pandurangan, Changho Choi, Vijay Balakrishnan |
IEEE CLOUD | 2 |
| 2018 | owlBIT: Orchestrating Wireless Transmissions for Launching Big Data Platforms in an Internet of Things EnvironmentabstractThe emergence of Edge Computing and the success of Internet of Thing and (IoT) has tremendously changed the way we think about data computing. With edge devices changing from data producer to both data producer and consumer, the chance for processing large data sets with Big Data on a cloud of IoT devices is more realistic. In Big Data systems such as Hadoop-Yarn, Spark, Pig, etc., the shuffling stage is by far the most dominant source of network traffic. Unreliable performance of network will greatly impact the shuffling process. Since IoT devices mostly rely on wireless network based on 802.11, providing proper throughput to big data computing system is an important challenge that needs to be address. In this paper, we argue that a cluster of IoT computers can support big data by considering the information fed by the big data applications. We propose a cross-layer framework that uses the application layer information to guide the packet scheduling at the link layer. We implement our system as an extension module in Hadoop-Yarn system. The experimental evaluation shows significant performance improvement. Son Nam Nguyen, Tengpeng Li, Bo Sheng, Ningfang Mi |
IEEE CLOUD | 6 |
| 2018 | Intermediate Data Caching Optimization for Multi-Stage and Parallel Big Data FrameworksabstractIn the era of big data and cloud computing, large amounts of data are generated from user applications and need to be processed in the datacenter. Data-parallel computing frameworks, such as Apache Spark, are widely used to perform such data processing at scale. Specifically, Spark leverages distributed memory to cache the intermediate results, represented as Resilient Distributed Datasets (RDDs). This gives Spark an advantage over other parallel frameworks for implementations of iterative machine learning and data mining algorithms, by avoiding repeated computation or hard disk accesses to retrieve RDDs. By default, caching decisions are left at the programmer's discretion, and the LRU policy is used for evicting RDDs when the cache is full. However, when the objective is to minimize total work, LRU is woefully inadequate, leading to arbitrarily suboptimal caching decisions. In this paper, we design an algorithm for multi-stage big data processing platforms to adaptively determine and cache the most valuable intermediate datasets that can be reused in the future. Our solution automates the decision of which RDDs to cache: this amounts to identifying nodes in a direct acyclic graph (DAG) representing computations whose outputs should persist in the memory. Our experiment results show that our proposed cache optimization solution can improve the performance of machine learning applications on Spark decreasing the total work to recompute RDDs by 12%. Zhengyu Yang 0001, Danlin Jia, Stratis Ioannidis, Ningfang Mi, Bo Sheng |
IEEE CLOUD | 4 |
| 2018 | GEDetector: Early Detection of Gathering Events Based on Cluster Containment Join in Trajectory Streams
Bin Zhao 0002, Genlin Ji, Zhaoyuan Yu, Xintao Liu, Ningfang Mi |
EDBT | 6 |
| 2018 | RoVEr: Robust and Verifiable Erasure Code for Hadoop Distributed File SystemsabstractErasure Coding based Storage (ECS) is replacing tradition replica-based systems because of its low storage overhead. In an ECS, however, every task needs to fetch remote pieces of data for its execution, and data verification is missing in the current framework. As security issues keep rising and there have been security incidents occurred in big data platforms, the compromised nodes in a computing cluster may manipulate its hosted data fed for other nodes yielding misleading results. Without replicas, it is quite challenging to efficiently verify the data integrity in ECS. In this paper, we develop ROVER, which is an efficient and verifiable ECS for big data platforms. In ROVER, every piece of data is monitored by its checksums stored on a set of witnesses. Bloom filter technique is used on each witness to efficiently keep the records of the checksums. The data verification is based on the majority voting. ROVER also supports a quick reconstruction of Bloom Filter when a node recovers from a failure. We present a complete system framework, security analysis, and a guideline for setting the parameters. The implementation and evaluation show that ROVER is robust and efficient against the attack from the compromised nodes. Son Nam Nguyen, Tengpeng Li, Ningfang Mi, Bo Sheng |
ICCCN | 6 |
| 2017 | FIM: Performance Prediction for Parallel Computation in Iterative Data Processing ApplicationsabstractPredicting performance of an application running on high performance computing (HPC) platforms in a cloud environment is increasingly becoming important because of its influence on development time and resource management. However, predicting the performance with respect to parallel processes is complex for iterative, multi-stage applications. This research proposes a performance approximation approach FiM to model the computing performance of iterative, multi-stage applications running on a master-compute framework. FiM consists of two key components that are coupled with each other: 1) Stochastic Markov Model to capture non-deterministic runtime that often depends on parallel resources, e.g., number of processes. 2) Machine Learning Model that extrapolates the parameters for calibrating our Markov model when we have changes in application parameters such as dataset. Our new modeling approach considers different design choices along multiple dimensions, namely (i) process level parallelism, (ii) distribution of cores on multi-core processors in cloud computing, (iii) application related parameters, and (iv) characteristics of datasets. The major contribution of our prediction approach is that FiM is able to provide an accurate prediction of parallel computation time for the datasets which have much larger size than that of the training datasets. Such calculation prediction provides data analysts a useful insight of optimal configuration of parallel resources (e.g., number of processes and number of cores) and also helps system designers to investigate the impact of changes in application parameters on system performance. Janki Bhimani, Ningfang Mi, Miriam Leeser, Zhengyu Yang 0001 |
CLOUD | 2 |
| 2017 | A Case for Abstract Cost Models for Distributed Execution of Analytics Operators
Rundong Li 0001, Ningfang Mi, Mirek Riedewald, Yizhou Sun |
DaWaK | 2 |
| 2017 | AutoPath: Harnessing Parallel Execution Paths for Efficient Resource Allocation in Multi-Stage Big Data FrameworksabstractDue to the flexibility of data operations and scalability of in- memory cache, Spark has revealed the potential to become the standard distributed framework to replace Hadoop for data-intensive processing in both industry and academia. However, we observe that the built-in scheduling algorithms in Spark (i.e., FIFO and FAIR) are not optimized for the applications with multiple parallel and independent branches in stages. Specifically, the child stage needs to wait and collect data from all its parent branches, but this wait has no guaranteed upper bound since it is tightly coupled with each branch's workload characteristic, stage order, and their corresponding allocated computing resource. To address this challenge, we investigate a superior solution which ensures all branches acquire suitable resources according to their workload demand in order to let the finish time of each branch be as close as possible. Based on this, we propose a novel scheduling policy, named AutoPath, which can effectively reduce the overall makespan of such kind of applications by detecting and leveraging the parallel path, and adaptively assigning computing resources based on the estimated workload demands during runtime. We implemented the new scheduling scheme in Spark v1.5.0 and evaluated it with selected representative workloads. The experiments demonstrate that our new scheduler effectively reduces the makespan and improves resource utilizations for these applications, compared to the current FIFO and FAIR schedulers. Han Gao 0013, Zhengyu Yang 0001, Janki Bhimani, Bo Sheng, Ningfang Mi |
ICCCN | 7 |
| 2017 | EA2S2: An Efficient Application-Aware Storage System for Big Data Processing in Heterogeneous ClustersabstractBig data processing frameworks such as Hadoop have been widely adopted to process a large volume of data. A lot of prior work has focused on the allocation of resources and the execution order of jobs/tasks to improve the performance in a homogeneous cluster. In this paper, we investigate storage layer design in a heterogeneous system considering a new type of bundled jobs where the input data and associated application jobs are submitted in a bundle. Our goal is to break the barrier between resource management and the underlying storage layer, and improve data locality, an important performance factor for resource management, from the aspect of storage system. We develop a sampling-based randomized algorithm for the network file system to determine the placement of input data blocks. The main idea is to query a selected set of candidate nodes, and estimate their workload at run time combining centralized and per-node information. The node with the smallest workload is selected to host the data block. Our evaluation is based with system implementation and comprehensive experiments on NSF CloudLab platforms. We have also conducted simulation for large-scale clusters. The results show significant performance improvements in terms of execution time and data locality. Son Nam Nguyen, Zhengyu Yang 0001, Ningfang Mi, Bo Sheng |
ICCCN | 5 |
| 2017 | Enhancing SSDs with multi-stream: What? why? how?abstractThe adoption of SSDs has become very prominent, but they still suffer from challenges to control write amplification. Traditional SSDs have single active append point where new data writes can be stored. Data of different lifetime stored together causes high write amplification. Recently, multi-stream SSDs are developed that allows multiple active append points. These multiple active append points can be used to store data of different lifetime in different locations within SSD. Such a data placement according to the lifetime of data would considerably reduce internal write amplification of SSD. For using multi-stream SSDs it is required to attach stream-id to each new incoming data writes. According to these stream-ids, the flash transition layer (FTL) of a multi-stream SSD then appends data to different erase blocks. Thus, multi-stream SSDs will help to reduce write amplification. But, to efficiently use this new multi-stream SSDs, it is important to properly identify streamids of data with respect to its lifetime. The lifetime of data is expected using different features that data exhibits like frequency, sequentiality etc. Stream-id identification using different features may have different impact on the final write amplification of multi-stream SSDs, depending on workload. Thus, it is required to quantify the impact of different data features that are used for stream-id identification on the resultant write amplification. Additionally, the combination of these data features may be used for stream-id identification, so it is also important to be able to study the impact of such different combinations. In order to address above challenges towards efficiently using multi-stream SSDs, here we propose a portable and adoptable framework to study the impact of stream-id identification using different workload data features and their combinations on write amplification of multi-stream SSDs. Our evaluation results show that use of appropriate features according to workload can considerably reduce the Write Amplification Factor (WAF) when compared to the legacy SSDs. Janki Bhimani, Jingpei Yang, Zhengyu Yang 0001, Ningfang Mi, N. H. V. Krishna Giri, Rajinikanth Pandurangan, Changho Choi, Vijay Balakrishnan |
IPCCC | 4 |
| 2017 | Cyber-physical system enabled nearby traffic flow modelling for autonomous vehiclesabstractWe propose a nearby traffic flow modelling solution based on built-in Cyber-Physical System (CPS) sensors of autonomous vehicles. Our goal is to enhance the offline route planning and driving decision adjustment based on the first-hand traffic information, especially during poor Internet connection moments. Specifically, our model helps to select the optimal speed on a road, the optimal distance for timing to brake, and the safe distance from other vehicles to keep. Moreover, our model can also assist neighboring autonomous vehicles by communicating required information through Ad-Hoc network communications or through a centralized cloud. In detail, we first focus on the unique characteristic of traffic flow (such as traffic rule, avoid collision behaviours), and then build a comprehensive model to handle multiple scenarios. Technically, our model uses density functions of velocities, the differential equation of traffic flows, and the traffic viscosity with information collected from the traffic flow, the distances between vehicles, the amount and density of vehicle, the instant velocity, the speed limit, and the momentum to analysis the the driving scene. We evaluate our model with real traffic data collected by in-vehicle CPS sensors to the proposed nearby traffic flow model. Results show that our work can accurately conduct offline estimation on nearby traffic signal influence, and reveal the correlations among velocity, density and (spatial and temporal) location to adjust route during runtime. Zhengyu Yang 0001, Siyu Huang, Xianzhi Du, Janki Bhimani, Ningfang Mi |
IPCCC | 8 |
| 2017 | AutoTiering: Automatic data placement manager in multi-tier all-flash datacenterabstractIn the year of 2017, the capital expenditure of Flash-based Solid State Drivers (SSDs) keeps declining and the storage capacity of SSDs keeps increasing. As a result, the “selling point” of traditional spinning Hard Disk Drives (HDDs) as a backend storage — low cost and large capacity — is no longer unique, and eventually they will be replaced by low-end SSDs which have large capacity but perform orders of magnitude better than HDDs. Thus, it is widely believed that all-flash multi-tier storage systems will be adopted in the enterprise datacenters in the near future. However, existing caching or tiering solutions for SSD-HDD hybrid storage systems are not suitable for all-flash storage systems. This is because that all-flash storage systems do not have a large speed difference (e.g., 10x) among each tier. Instead, different specialties (such as high performance, high capacity, etc.) of each tier should be taken into consideration. Motivated by this, we develop an automatic data placement manager called “AutoTiering” to handle virtual machine disk files (VMDK) allocation and migration in an all-flash multitier datacenter to best utilize the storage resource, optimize the performance, and reduce the migration overhead. AutoTiering is based on an optimization framework, whose core technique is to predict VM's performance change on different tiers with different specialties without conducting real migration. As far as we know, AutoTiering is the first optimization solution designed for all-flash multi-tier datacenters. We implement AutoTiering on VMware ESXi [1], and experimental results show that it can significantly improve the I/O performance compared to existing solutions. Zhengyu Yang 0001, Morteza Hoseinzadeh, Allen Andrews, Clay Mayers, David Thomas Evans, Rory Thomas Bolt, Janki Bhimani, Ningfang Mi, Steven Swanson |
IPCCC | 8 |
| 2017 | H-NVMe: A hybrid framework of NVMe-based storage system in cloud computing environmentabstractIn the year of 2017, more and more datacenters have started to replace traditional SATA and SAS SSDs with NVMe SSDs due to NVMe's outstanding performance [1]. However, for historical reasons, current popular deployments of NVMe in VM-hypervisor-based platforms (such as VMware ESXi [2]) have numbers of intermediate queues along the I/O stack. As a result, performance is bottlenecked by synchronization locks in these queues, cross-VM interference induces I/O latency, and most importantly, up-to-64K-queue capability of NVMe SSDs cannot be fully utilized. In this paper, we developed a hybrid framework of NVMe-based storage system called “H-NVMe”, which provides two VM I/O stack deployment modes “Parallel Queue Mode” and “Direct Access Mode”. The first mode increases parallelism and enables lock-free operations by implementing local lightweight queues in the NVMe driver. The second mode further bypasses the entire I/O stack in the hypervisor layer and allows trusted user applications whose hosting VMDK (Virtual Machine Disk) files are attached with our customized vSphere IOFilters [3] to directly access NVMe SSDs to improve the performance isolation. This suits premium users who have higher priorities and the permission to attach IOFilter to their VMDKs. H-NVMe is implemented on VMware EXSi 6.0.0, and our evaluation results show that the proposed H-NVMe framework can significant improve throughputs and bandwidths compared to the original inbox NVMe solution. Zhengyu Yang 0001, Morteza Hoseinzadeh, Ping Wong, John Artoux, Clay Mayers, David Thomas Evans, Rory Thomas Bolt, Janki Bhimani, Ningfang Mi, Steven Swanson |
IPCCC | 9 |
| 2017 | Improving Flash Resource Utilization at Minimal Management Cost in Virtualized Flash-Based Storage SystemsabstractEffectively leveraging Flash resources has emerged as a highly important problem in enterprise storage systems. One of the popular techniques today is to use Flash as a secondary-level host-side cache in the virtual machine environment. Although this approach delivers IO acceleration for VMs' IO workloads, it might not be able to fully exploit the outstanding performance of Flash and justify the high cost-per-GB of Flash resources. In this paper, we design new VMware Flash Resource Managers (VFRM and GLB-VFRM) under the consideration of both performance and the incurred cost for managing Flash resources. Specifically, VFRM and GLB-VFRM aim to maximize the utilization of Flash resources with minimal CPU, memory and IO cost in managing and operating Flash for a dedicated enterprise workload and multiple heterogeneous enterprise workloads, respectively. Our new Flash resource managers adopt the ideas of thermodynamic heating and cooling to identify data blocks that can benefit the most from being put on Flash and migrate data blocks between Flash and magnetic disks in a lazy and asynchronous mode. Experimental evaluation of the prototype shows that both VFRM and GLB-VFRM achieve better cost-effectiveness than traditional caching solutions, i.e., obtaining IO hit ratios even slightly better than some of the conventional algorithms as Flash size increases yet costing orders of magnitude less IO bandwidth. Jianzhe Tai, Deng Liu, Zhengyu Yang 0001, Xiaoyun Zhu, Jack Lo, Ningfang Mi |
IEEE Trans. Cloud Comput. | 6 |
| 2017 | Self-Adjusting Slot Configurations for Homogeneous and Heterogeneous Hadoop ClustersabstractThe MapReduce framework and its open source implementation Hadoop have become the defacto platform for scalable analysis on large data sets in recent years. One of the primary concerns in Hadoop is how to minimize the completion length (i.e., makespan) of a set of MapReduce jobs. The current Hadoop only allows static slot configuration, i.e., fixed numbers of map slots and reduce slots throughout the lifetime of a cluster. However, we found that such a static configuration may lead to low system resource utilizations as well as long completion length. Motivated by this, we propose simple yet effective schemes which use slot ratio between map and reduce tasks as a tunable knob for reducing the makespan of a given set. By leveraging the workload information of recently completed jobs, our schemes dynamically allocates resources (or slots) to map and reduce tasks. We implemented the presented schemes in Hadoop V0.20.2 and evaluated them with representative MapReduce benchmarks at Amazon EC2. The experimental results demonstrate the effectiveness and robustness of our schemes under both simple workloads and more complex mixed workloads. Bo Sheng, Chiu C. Tan 0001, Ningfang Mi |
IEEE Trans. Cloud Comput. | 5 |
| 2016 | A Fresh Perspective on Total Cost of Ownership Models for Flash Storage in DatacentersabstractRecently, adoption of Flash based devices has become increasingly common in all forms of computing devices. Flash devices have started to become more economically viable for large storage installations like datacenters, where metrics like Total Cost of Ownership (TCO) are of paramount importance. Flash devices suffer from write amplification (WA), which, if unaccounted, can substantially increase the TCO of a storage system. In this paper, we develop a TCO model for Flash storage devices, and then plug a Write Amplification (WA) model of NVMe SSDs we build based on empirical data into this TCO model. Our new WA model accounts for workload characteristics like write rate and percentage of sequential writes. Furthermore, using both the TCO and WA models as the optimization criterion, we design new Flash resource management schemes (minTCO) to guide datacenter managers to make workload allocation decisions under the consideration of TCO for SSDs. Experimental results show that minTCO can reduce the TCO and keep relatively high throughput and space utilization of the entire datacenter storage. Zhengyu Yang 0001, Manu Awasthi, Mrinmoy Ghosh, Ningfang Mi |
CloudCom | 4 |
| 2016 | OpERA: Opportunistic and Efficient Resource Allocation in Hadoop YARN by Harnessing Idle ResourcesabstractEfficiently managing resources and improving throughput in a large-scale cluster has become a crucial problem with the explosion of data processing applications in recent years. Hadoop YARN and Mesos, as two universal resource management platforms, have been widely adopted in the commodity cluster for co-deploying multiple data processing frameworks, such as Hadoop MapReduce and Apache Spark. However, in the existing resource management, a certain amount of resources are exclusively allocated to a running task and can only be re-assigned after that task is completed. This exclusive mode unfortunately leads to a potential problem that may underutilize the cluster resources and degrade system performance. To address this issue, we propose a novel opportunistic and efficient resource allocation approach, named OpERA, which breaks the barriers among the encapsulated resource containers by leveraging the knowledge of actual runtime resource utilizations to re-assign opportunistic available resources to the pending tasks. We implement and evaluate OpERA in Hadoop YARN v2.5. Our experimental results show that OpERA significantly reduces the average job execution time and increases the resource (CPU and memory) utilizations. Ningfang Mi, Bo Sheng |
ICCCN | 4 |
| 2016 | Performance prediction techniques for scalable large data processing in distributed MPI systemsabstractPredicting performance of an application running on parallel computing platforms is increasingly becoming important due to the long development time of an application and the high resource management cost of parallel computing platforms. However, predicting overall performance is complex and must take into account both parallel calculation time and communication time. Difficulty in accurate performance modeling is compounded by myriad design choices along multiple dimensions, namely (i) process level parallelism, (ii) distribution of cores on multi-processor platforms, (iii) application related parameters, and (iv) characteristics of datasets. This research proposes a fast and accurate performance prediction approach to predict the calculation and communication time of an application running on a distributed computing platform. The major contribution of our prediction approach is that it can provide an accurate prediction of execution times for new datasets which have much larger sizes than the training datasets. Our approach consists of two models, i.e., a probabilistic self-learning model to predict calculation time and a simulation queuing model to predict network communication time. The combination of these two models provides data analysts a useful insight of optimal configuration of parallel resources (e.g., number of processes and number of cores) and application parameters setting. Janki Bhimani, Ningfang Mi, Miriam Leeser |
IPCCC | 2 |
| 2016 | Understanding performance of I/O intensive containerized applications for NVMe SSDsabstractOur cloud-based IT world is founded on hyper-visors and containers. Containers are becoming an important cornerstone, which is increasingly used day-by-day. Among different available frameworks, docker has become one of the major adoptees to use containerized platform in data centers and enterprise servers, due to its ease of deploying and scaling. Further more, the performance benefits of a lightweight container platform can be leveraged even more with a fast back-end storage like high performance SSDs. However, increase in number of simultaneously operating docker containers may not guarantee an aggregated performance improvement due to saturation. Thus, understanding performance bottleneck in a multi-tenancy docker environment is critically important to maintain application level fairness and perform better resource management. In this paper, we characterize the performance of persistent storage option (through data volume) for I/O intensive, dockerized applications. Our work investigates the impact on performance with increasing number of simultaneous docker containers in different workload environments. We provide, first of its kind study of I/O intensive containerized applications operating with NVMe SSDs. We show that 1) a six times better application throughput can be obtained, just by wise selection of number of containerized instances compared to single instance; and 2) for multiple application containers running simultaneously, an application throughput may degrade upto 50% compared to a stand-alone applications throughput, if good choice of application and workload is not made. We then propose novel design guidelines for an optimal and fair operation of both homogeneous and heterogeneous environments mixed with different applications and workloads. Janki Bhimani, Jingpei Yang, Zhengyu Yang 0001, Ningfang Mi, Qiumin Xu, Manu Awasthi, Rajinikanth Pandurangan, Vijay Balakrishnan |
IPCCC | 4 |
| 2016 | eSplash: Efficient speculation in large scale heterogeneous computing systemsabstractIn this paper, we aim to develop an efficient speculation framework for a heterogeneous cluster. Speculation is a common mechanism that identifies ‘slow’ node in a cluster and starts redundant tasks on other nodes to guarantee the reliability. We consider MapReduce/Hadoop as a representative computing platform, and our general goal is to accurately and quickly identify the straggler nodes during the job execution. On the one hand, our approach significantly reduces unnecessary speculative executions that occupy system resources, but do not get finished. On the other hand, when a node is prone to failure, our solution is able to detect it at an early stage and effectively launch a speculative task to avoid the delay in the job execution. We implement our solution in Hadoop platform and evaluate it with extensive experiments. The results show that our solution is efficient and effective when handling the speculative execution. The job execution time in our system is superior to that in the current Hadoop distribution. Zhengyu Yang 0001, Ningfang Mi, Bo Sheng |
IPCCC | 4 |
| 2016 | GReM: Dynamic SSD resource allocation in virtualized storage systems with heterogeneous IO workloadsabstractIn a shared virtualized storage system that runs VMs with heterogeneous IO demands, it becomes a problem for the hypervisor to cost-effectively partition and allocate SSD resources among multiple VMs. There are two straightforward approaches to solving this problem: equally assigning SSDs to each VM or managing SSD resources in a fair competition mode. Unfortunately, neither of these approaches can fully utilize the benefits of SSD resources, particularly when the workloads frequently change and bursty IOs occur from time to time. In this paper, we design a Global SSD Resource Management solution - GReM, which aims to fully utilize SSD resources as a second-level cache under the consideration of performance isolation. In particular, GReM takes dynamic IO demands of all VMs into consideration to split the entire SSD space into a long-term zone and a short-term zone, and cost-effectively updates the content of SSDs in these two zones. GReM is able to adaptively adjust the reservation for each VM inside the long-term zone based on their IO changes. GReM can further dynamically partition SSDs between the long- and short-term zones during runtime by leveraging the feedbacks from both cache performance and bursty workloads. Experimental results show that GReM can capture the cross-VM IO changes to make correct decisions on resource allocation, and thus obtain high IO hit ratio and low IO management costs, compared with both traditional and state-of-the-art caching algorithms. Zhengyu Yang 0001, Jianzhe Tai, Janki Bhimani, Ningfang Mi, Bo Sheng |
IPCCC | 5 |
| 2016 | AutoReplica: Automatic data replica manager in distributed caching and data processing systemsabstractNowadays, replication technique is widely used in data center storage systems for large scale Cyber-physical Systems (CPS) to prevent data loss. However, side-effect of replication is mainly the overhead of extra network and I/O traffics, which inevitably downgrades the overall I/O performance of the cluster. To effectively balance the trade-off between I/O performance and fault tolerance, in this paper, we propose a complete solution called “AutoReplica” - a replica manager in distributed caching and data processing systems with SSD-HDD tier storages. In detail, AutoReplica utilizes the remote SSDs (connected by high speed fibers) to replicate local SSD caches to protect data. In order to conduct load balancing among nodes and reduce the network overhead, we propose three approaches (i.e., ring, network, and multiple-SLA network) to automatically setup the cross-node replica structure with the consideration of network traffic, I/O speed and SLAs. To improve the performance during migrations triggered by load balance and failure recovery, we propose the a migrate-on-write technique called “fusion cache” to seamlessly migrate and prefetch among local and remote replicas without pausing the subsystem. Moreover, AutoReplica can also recover from different failure scenarios, while limits the performance downgrading degree. Lastly, AutoReplica supports parallel prefetching from multiple nodes with a new dynamic optimizing streaming technique to improve I/O performance. We are currently in the process of implementing AutoReplica to be easily plugged into commonly used distributed caching systems, and solidifying our design and implementation details. Zhengyu Yang 0001, David Thomas Evans, Ningfang Mi |
IPCCC | 4 |
| 2016 | A new packet scheduling algorithm for access points in crowded WLANs
Bo Sheng, Ningfang Mi |
Ad Hoc Networks | 3 |
| 2015 | Admission control in YARN clusters based on dynamic resource reservationabstractHadoop YARN is an open project developed by the Apache Software Foundation to provide a resource management framework for large scale parallel data processing. However, there exists a resource waiting deadlock under the Fair scheduler when the resource requisition of applications is beyond the amount that the cluster can provide. In such a case, the YARN system will be halted if all resources are occupied by ApplicationMasters, a special task of each job that negotiates resources for processing tasks and coordinates job execution. Therefore, we develop a new admission control mechanism which dynamically reserves resources for processing tasks in order to avoid resource waiting deadlocks and meanwhile obtain good performance. We implement and evaluate our new mechanism in Hadoop YARN v2.2.0. The experimental results show the effectiveness of this mechanism under MapReduce benchmarks. Ningfang Mi, Bo Sheng |
IM | 4 |
| 2015 | OMO: Optimize MapReduce overlap with a good start (reduce) and a good finish (map)abstractMapReduce has become a popular data processing framework in the past few years. Scheduling algorithm is crucial to the performance of a MapReduce cluster, especially when the cluster is concurrently executing a batch of MapReduce jobs. However, the scheduling problem in MapReduce is different from the traditional job scheduling problem as the reduce phase usually starts before the map phase is finished to “shuffle” the intermediate data. This paper develops a new strategy, named OMO, which particularly aims to optimize the overlap between the map and reduce phases. Our solution includes two new techniques, lazy start of reduce tasks and batch finish of map tasks, which catch the characteristics of the overlap in a MapReduce process and achieve a good alignment of the two phases. We have implemented OMO on Hadoop system and evaluated the performance with extensive experiments. The results show that OMO's performance is superior in terms of total completion length (i.e., makespan) of a batch of jobs. Ying Mao 0001, Bo Sheng, Ningfang Mi |
IPCCC | 5 |
| 2015 | LsPS: A Job Size-Based Scheduler for Efficient Task Assignments in HadoopabstractThe MapReduce paradigm and its open source implementation Hadoop are emerging as an important standard for large-scale data-intensive processing in both industry and academia. A MapReduce cluster is typically shared among multiple users with different types of workloads. When a flock of jobs are concurrently submitted to a MapReduce cluster, they compete for the shared resources and the overall system performance in terms of job response times, might be seriously degraded. Therefore, one challenging issue is the ability of efficient scheduling in such a shared MapReduce environment. However, we find that conventional scheduling algorithms supported by Hadoop cannot always guarantee good average response times under different workloads. To address this issue, we propose a new Hadoop scheduler, which leverages the knowledge of workload patterns to reduce average job response times by dynamically tuning the resource shares among users and the scheduling algorithms for each user. Both simulation and real experimental results from Amazon EC2 cluster show that our scheduler reduces the average MapReduce job response time under a variety of system workloads compared to the existing FIFO and Fair schedulers. Jianzhe Tai, Bo Sheng, Ningfang Mi |
IEEE Trans. Cloud Comput. | 4 |
| 2014 | Using Elasticity to Improve Inline Data Deduplication Storage SystemsabstractElasticity is the ability to scale computing resources such as memory on-demand, and is one of the main advantages of utilizing cloud computing services. With the increasing popularity of cloud based storage, it is natural that more deduplication based storage systems will be migrated to the cloud. Existing deduplication systems however, do not adequately take advantage of elasticity. In this paper, we illustrate how to use elasticity to improve deduplication based systems, and propose EAD (elasticity aware deduplication), an indexing algorithm that uses the ability to dynamically increase memory resources to improve overall deduplication performance. Our experimental results indicate that EAD is able to detect more than 98\% of all duplicate data, however only consumes less than 5\% of expected memory space. Meanwhile, it claims four times of deduplication efficiency than the state-of-art sampling technique while costs less than half of the amount of memory. Yufeng Wang 0007, Chiu C. Tan 0001, Ningfang Mi |
IEEE CLOUD | 3 |
| 2014 | FRESH: Fair and Efficient Slot Configuration and Scheduling for Hadoop ClustersabstractHadoop is an emerging framework for parallel big data processing. While becoming popular, Hadoop is too complex for regular users to fully understand all the system parameters and tune them appropriately. Especially when processing a batch of jobs, default Hadoop setting may cause inefficient resource utilization and unnecessarily prolong the execution time. This paper considers an extremely important setting of slot configuration which by default is fixed and static. We proposed an enhanced Hadoop system called FRESH which can derive the best slot setting, dynamically configure slots, and appropriately assign tasks to the available slots. The experimental results show that when serving a batch of MapReduce jobs, FRESH significantly improves the makespan as well as the fairness among jobs. Ying Mao 0001, Bo Sheng, Ningfang Mi |
IEEE CLOUD | 5 |
| 2014 | HaSTE: Hadoop YARN Scheduling Based on Task-Dependency and Resource-DemandabstractThe MapReduce framework has become the de facto scheme for scalable semi-structured and un-structured data processing in recent years. The Hadoop ecosystem has evolved into its second generation, Hadoop YARN, which adopts fine-grained resource management schemes for job scheduling. One of the primary performance concerns in YARN is how to minimize the total completion length, i.e., makespan, of a set of MapReduce jobs. However, the precedence constraint or fairness constraint in current widely used scheduling policies in YARN, such as FIFO and Fair, can both lead to inefficient resource allocation in the Hadoop YARN cluster. They also omit the dependency between tasks which is crucial for the efficiency of resource utilization. We thus propose a new YARN scheduler, named HaSTE, which can effectively reduce the makespan of MapReduce jobs in YARN by leveraging the information of requested resources, resource capacities, and dependency between tasks. We implemented HaSTE as a pluggable scheduler in the most recent version of Hadoop YARN, and evaluated it with classic MapReduce benchmarks. The experimental results demonstrate that our YARN scheduler effectively reduces the makespans and improves resource utilization compare to the current scheduling policies. Bo Sheng, Ningfang Mi |
IEEE CLOUD | 5 |
| 2014 | Live Data Migration for Reducing SLA Violations in Multi-tiered Storage SystemsabstractToday, the volume of data in the world has been tremendously increased. Large-scaled and diverse data sets are raising new big challenges of storage, process, and query. Tiered storage architectures combining solid-state drives (SSDs) with hard disk drives (HDDs), become attractive in enterprise data centers for achieving high performance and large capacity simultaneously. However, how to best use these storage resources and efficiently manage massive data for providing high quality of service (QoS) is still a core and difficult problem. In this paper, we present a new approach for automated data movement in multi-tiered storage systems, which lively migrates the data across different tiers, aiming to support multiple service level agreements (SLAs) for applications with dynamic workloads at the minimal cost. Trace-driven simulations show that compared to the no migration policy, LMsT significantly improves average I/O response times, I/O violation ratios and I/O violation times, with only slight degradation (e.g., up to 6% increase in SLA violation ratio) on the performance of high priority applications. Jianzhe Tai, Bo Sheng, Ningfang Mi |
IC2E | 4 |
| 2014 | Improving Virtual Machine Migration via DeduplicationabstractFor this study the techniques of virtual machine migration are understood and the affects deduplication has on migration are evaluated. The benefits of using deduplication and compression on virtual machines show in the metric of space saved during migrating. Deduplication is computationally expensive so we evaluate how to group virtual machines with similar elements in order to improve migration. From this study, grouping virtual machines based on similar elements improves the overhead from deduplication and compression but estimates which virtual machines are best grouped together. Jake Roemer, Mark Groman, Zhengyu Yang 0001, Yufeng Wang 0007, Chiu C. Tan 0001, Ningfang Mi |
MASS | 6 |
| 2014 | VFRM: Flash Resource Manager in VMware ESX ServerabstractOne popular approach of leveraging Flash technology in the virtual machine environment today is using it as a secondary-level host-side cache. Although this approach delivers I/O acceleration for a single VM workload, it might not be able to fully exploit the outstanding performance of Flash and justify the high cost-per-GB of Flash resources. In this paper, we present the design for VMware Flash Resource Manager (VFRM), which aims to maximize the utilization of Flash resources with minimal CPU, memory and I/O cost for managing and operating Flash. It borrows the ideas of heating and cooling from thermodynamics to identify the data blocks that benefit most from being put on Flash, and lazily and asynchronously migrates the data blocks between Flash and spinning disks. Experimental evaluation of the prototype shows that VFRM achieves better cost-effectiveness than traditional caching solutions, and costs orders of magnitude less memory and I/O bandwidth. Deng Liu, Ningfang Mi, Jianzhe Tai, Xiaoyun Zhu, Jack Lo |
NOMS | 2 |
| 2013 | Using a Tunable Knob for Reducing Makespan of MapReduce Jobs in a Hadoop ClusterabstractThe MapReduce framework and its open source implementation Hadoop have become the defacto platform for scalable analysis on large data sets in recent years. One of the primary concerns in Hadoop is how to minimize the completion length (i.e., makespan) of a set of MapReduce jobs. The current Hadoop only allows static slot configuration, i.e., fixed numbers of map slots and reduce slots throughout the lifetime of a cluster. However, we found that such a static configuration may lead to low system resource utilizations as well as long completion length. Motivated by this, we propose a simple yet effective scheme which uses slot ratio between map and reduce tasks as a tunable knob for reducing the makespan of a given set. By leveraging the workload information of recently completed jobs, our scheme dynamically allocates resources (or slots) to map and reduce tasks. We implemented the presented scheme in Hadoop V0.20.2 and evaluated it with representative MapReduce benchmarks at Amazon EC2. The experimental results demonstrate the effectiveness and robustness of our scheme under both simple workloads and more complex mixed workloads. Bo Sheng, Ningfang Mi |
IEEE CLOUD | 4 |
| 2013 | Scheduling heterogeneous MapReduce jobs for efficiency improvement in enterprise clusters
Jianzhe Tai, Bo Sheng, Ningfang Mi |
IM | 4 |
| 2012 | ADuS: Adaptive resource allocation in cluster systems under heavy-tailed and bursty workloadsabstractA large-scaled cluster system has been employed in various areas by offering pools of fundamental resources. How to effectively allocate the shared resources in a cluster system is a critical but challenging issue, which has been extensively studied in the past few years. Despite the fact that classic load balancing policies, such as Random, Join Shortest Queue and size-based polices, are widely implemented in actual systems due to their simplicity and efficiency, the performance benefits of these policies diminish when workloads are highly variable and heavily dependent. In this paper, we propose a new load balancing policy named ADuS, which attempts to partition jobs according to their sizes and to further rank the servers based on their loads. By dispatching jobs of similar size to the servers with the same ranking, ADuS can adaptively balance user traffic and system load in the system and thus achieve significant performance benefits. Extensive simulations show the effectiveness and the robustness of ADuS under many different environments. Jianzhe Tai, Ningfang Mi |
ICC | 4 |
| 2012 | ASIdE: Using Autocorrelation-Based Size Estimation for Scheduling Bursty WorkloadsabstractTemporal dependence in workloads creates peak congestion that can make service unavailable and reduce system performance. To improve system performability under conditions of temporal dependence, a server should quickly process bursts of requests that may need large service demands. In this paper, we propose and evaluateASIdE, an Autocorrelation-based SIze Estimation, that selectively delays requests which contribute to the workload temporal dependence. ASIdE implicitly approximates the shortest job first (SJF) scheduling policy but without any prior knowledge of job service times. Extensive experiments show that (1) ASIdE achieves good service time estimates from the temporal dependence structure of the workload to implicitly approximate the behavior of SJF; and (2) ASIdE successfully counteracts peak congestion in the workload and improves system performability under a wide variety of settings. Specifically, we show that system capacity under ASIdE is largely increased compared to the first-come first-served (FCFS) scheduling policy and is highly-competitive with SJF. Ningfang Mi, Giuliano Casale, Evgenia Smirni |
IEEE Trans. Netw. Serv. Manag. | 1 |
| 2012 | Dealing with Burstiness in Multi-Tier Applications: Models and Their ParameterizationabstractWorkloads and resource usage patterns in enterprise applications often show burstiness resulting in large degradation of the perceived user performance. In this paper, we propose a methodology for detecting burstiness symptoms in multi-tier applications but, rather than identifying the root cause of burstiness, we incorporate this information into models for performance prediction. The modeling methodology is based on the index of dispersion of the service process at a server, which is inferred by observing the number of completions within the concatenated busy times of that server. The index of dispersion is used to derive a Markov-modulated process that captures burstiness and variability of the service process at each resource well and that allows us to define queueing network models for performance prediction. Experimental results and performance model predictions are in excellent agreement and argue for the effectiveness of the proposed methodology under both bursty and nonbursty workloads. Furthermore, we show that the methodology extends to modeling flash crowds that create burstiness in the stream of requests incoming to the application. Giuliano Casale, Ningfang Mi, Ludmila Cherkasova, Evgenia Smirni |
IEEE Trans. Software Eng. | 2 |
| 2011 | Decentralized Scheduling of Bursty Workload on Computing GridsabstractBursty workloads are often observed in a variety of systems such as grid services, multi-tier architectures, and large storage systems. Studies have shown that such burstiness can dramatically degrade system performance because of overloading, increased response time, and unavailable service. Computing grids, which often use distributed, autonomous resource management, are particularly susceptible to load imbalances caused by bursty workloads. In this paper, we use a simulation environment to investigate the performance of decentralized schedulers under various intensity levels of burstiness. We first demonstrate a significant performance degradation in the presence of strong and moderate bursty workloads. Then, we describe two new hybrid schedulers, based on duplication-invalidation, and assess the effectiveness of these schedulers under different intensities of burstiness. Our simulation results show that compared to the conventional decentralized methods, the proposed schedulers achieve a 40% performance improvement under the bursty condition while obtaining similar performance in non-bursty conditions. Juemin Zhang, Ningfang Mi, Jianzhe Tai, Waleed Meleis |
ICC | 2 |
| 2011 | ArA: Adaptive resource allocation for cloud computing environments under bursty workloadsabstractCloud computing nowadays becomes quite popular among a community of cloud users by offering a variety of resources. However, burstiness in user demands often dramatically degrades the application performance. In order to satisfy peak user demands and meet Service Level Agreement (SLA), efficient resource allocation schemes are highly demanded in the cloud. However, we find that conventional load balancers unfortunately neglect cases of bursty arrivals and thus experience significant performance degradation. Motivated by this problem, we propose new burstiness-aware algorithms to balance bursty workloads across all computing sites, and thus to improve overall system performance. We present a smart load balancer, which leverages the knowledge of burstiness to predict the changes in user demands and on-the-fly shifts between the schemes that are “greedy” (i.e., always select the best site) and “random” (i.e., randomly select one) based on the predicted information. Both simulation and real experimental results show that this new load balancer can adapt quickly to the changes in user demands and thus improve performance by making a smart site selection for cloud users under both bursty and non-bursty workloads. Jianzhe Tai, Juemin Zhang, Waleed Meleis, Ningfang Mi |
IPCCC | 5 |
| 2011 | DAT: An AP scheduler using dynamically adjusted time windows for crowded WLANsabstractThis paper proposes a new packet scheduling algorithm for access points in a crowded 802.11 WLAN. Our goal is to improve the performance of efficiency (measured by packet response time or throughput) and fairness which often conflict with each other. Our solution aggregates both metrics and leverages the balance between them. The basic idea is to let the AP allocate different time windows for serving each client. According to the observed traffic, our algorithm dynamically shifts the weight between efficiency and fairness and strikes to improve the preferred metric without excessively degrading the other one. A valid queuing model is developed to evaluate the new algorithm's performance. Using trace-driven simulations, we show that our algorithm successfully balances the trade off between the efficiency and the fairness in a busy WLAN. Bo Sheng, Ningfang Mi |
IPCCC | 3 |
| 2010 | AWAIT: Efficient Overload Management for Busy Multi-tier Web Services under Bursty Workloads
Ludmila Cherkasova, Vittoria de Nitto Persone, Ningfang Mi, Evgenia Smirni |
ICWE | 4 |
| 2010 | CWS: a model-driven scheduling policy for correlated workloadsabstractWe define CWS, a non-preemptive scheduling policy for workloads with correlated job sizes. CWS tackles the scheduling problem by inferring the expected sizes of upcoming jobs based on the structure of correlations and on the outcome of past scheduling decisions. Size prediction is achieved using a class of Hidden Markov Models (HMM) with continuous observation densities that describe job sizes. We show how the forward-backward algorithm of HMMs applies effectively in scheduling applications and how it can be used to derive closed-form expressions for size prediction. This is particularly simple to implement in the case of observation densities that are phase-type (PH-type) distributed, where existing fitting methods for Markovian point processes may also simplify the parameterization of the HMM workload model.Based on the job size predictions, CWS emulates size-based policies which favor short jobs, with accuracy depending mainly on the HMM used to parametrize the scheduling algorithm. Extensive simulation and analysis illustrate that CWS is competitive with policies that assume exact information about the workload. Giuliano Casale, Ningfang Mi, Evgenia Smirni |
SIGMETRICS | 2 |
| 2010 | Efficient resource allocation and power saving in multi-tiered systemsabstractIn this paper, we present Fastrack, a parameter-free algorithm for dynamic resource provisioning that uses simple statistics to promptly distill information about changes in workload burstiness. This information, coupled with the application's end-to-end response times and system bottleneck characteristics, guide resource allocation that shows to be very effective under a broad variety of burstiness profiles and bottleneck scenarios. Andrew Caniff, Ningfang Mi, Ludmila Cherkasova, Evgenia Smirni |
WWW | 3 |
| 2010 | Model-Driven System Capacity Planning under Workload BurstinessabstractIn this paper, we define and study a new class of capacity planning models called MAP queueing networks. MAP queueing networks provide the first analytical methodology to describe and predict accurately the performance of complex systems operating under bursty workloads, such as multitier architectures or storage arrays. Burstiness is a feature that significantly degrades system performance and that cannot be captured explicitly by existing capacity planning models. MAP queueing networks address this limitation by describing computer systems as closed networks of servers whose service times are Markovian Arrival Processes (MAPs), a class of Markov-modulated point processes that can model general distributions and burstiness. In this paper, we show that MAP queueing networks provide reliable performance predictions even if the service processes are bursty. We propose a methodology to solve MAP queueing networks by two state space transformations, which we call Linear Reduction (LR) and Quadratic Reduction (QR). These transformations dramatically decrease the number of states in the underlying Markov chain of the queueing network model. From these reduced state spaces, we obtain two classes of bounds on arbitrary performance indexes, e.g., throughput, response time, and utilizations. Numerical experiments show that LR and QR bounds achieve good accuracy. We also illustrate the high effectiveness of the LR and QR bounds in the performance analysis of a real multitier architecture subject to TPC-W workloads that are characterized as bursty. These results promote MAP queueing networks as a new class of robust capacity planning models. Giuliano Casale, Ningfang Mi, Evgenia Smirni |
IEEE Trans. Computers | 2 |
| 2009 | Autocorrelation-driven load control in distributed systemsabstractIn this paper, we propose a new approach for the development of load control policies in autonomic multitier systems. We control system load in a completely new way compared to existing policies: we leverage on the autocorrelation of service times and show that autocorrelation can be used to forecast future service requirements of requests and adaptively control system load. To the best of our knowledge, this is the first direct application of autocorrelation of service times to autonomic load control. We propose ALoC and D ALoC, two autocorrelation-driven policies that drop a percentage of the load in order to meet pre-defined quality-of-service levels in a distributed system. Both policies are easy to implement and rely on minimal assumptions. In particular, D ALoC is a fully no-knowledge measurement-based policy that self-adjusts its load control parameters based only on policy targets and on statistical information of requests served in the past. We illustrate the effectiveness of these new policies in a distributed multi-server setting via detailed trace driven simulations. We show that if these policies are employed in the server with a temporal dependent service process, then end-to-end response time, across all servers, reduces up to 80% by only dropping at most 13% of the incoming requests. Using real traces, we also show that, in the constrained case of being able to drop only from a portion of the incoming workload, our policy still improves request response time by up to 30%. Ningfang Mi, Giuliano Casale, Qi Zhang 0012, Alma Riska, Evgenia Smirni |
MASCOTS | 1 |
| 2009 | Automated anomaly detection and performance modeling of enterprise applicationsabstractAutomated tools for understanding application behavior and its changes during the application lifecycle are essential for many performance analysis and debugging tasks. Application performance issues have an immediate impact on customer experience and satisfaction. A sudden slowdown of enterprise-wide application can effect a large population of customers, lead to delayed projects, and ultimately can result in company financial loss. Significantly shortened time between new software releases further exacerbates the problem of thoroughly evaluating the performance of an updated application. Our thesis is that online performance modeling should be a part of routine application monitoring. Early, informative warnings on significant changes in application performance should help service providers to timely identify and prevent performance problems and their negative impact on the service. We propose a novel framework for automated anomaly detection and application change analysis. It is based on integration of two complementary techniques: (i) a regression-based transaction model that reflects a resource consumption model of the application, and (ii) an application performance signature that provides a compact model of runtime behavior of the application. The proposed integrated framework provides a simple and powerful solution for anomaly detection and analysis of essential performance changes in application behavior. An additional benefit of the proposed approach is its simplicity: It is not intrusive and is based on monitoring data that is typically available in enterprise production environments. The introduced solution further enables the automation of capacity planning and resource provisioning tasks of multitier applications in rapidly evolving IT environments. Ludmila Cherkasova, Kivanc M. Ozonat, Ningfang Mi, Julie Symons, Evgenia Smirni |
ACM Trans. Comput. Syst. | 3 |
| 2009 | Efficient management of idleness in storage systemsabstractVarious activities that intend to enhance performance, reliability, and availability of storage systems are scheduled with low priority and served during idle times. Under such conditions, idleness becomes a valuable “resource” that needs to be efficiently managed. A common approach in system design is to be nonwork conserving by “idle waiting”, that is, delay the scheduling of background jobs to avoid slowing down upcoming foreground tasks. In this article, we complement “idle waiting” with the “estimation” of background work to be served in every idle interval to effectively manage the trade-off between the performance of foreground and background tasks. As a result, the storage system is better utilized without compromising foreground performance. Our analysis shows that if idle times have low variability, then idle waiting is not necessary. Only if idle times are highly variable does idle waiting become necessary to minimize the impact of background activity on foreground performance. We further show that if there is burstiness in idle intervals, then it is possible to predict accurately the length of incoming idle intervals and use this information to serve more background jobs without affecting foreground performance. Ningfang Mi, Alma Riska, Qi Zhang 0012, Evgenia Smirni, Erik Riedel |
ACM Trans. Storage | 1 |
| 2008 | Anomaly? application change? or workload change? towards automated detection of application performance anomaly and changeabstractAutomated tools for understanding application behavior and its changes during the application life-cycle are essential for many performance analysis and debugging tasks. Application performance issues have an immediate impact on customer experience and satisfaction. A sudden slowdown of enterprise-wide application can effect a large population of customers, lead to delayed projects and ultimately can result in company financial loss. We believe that online performance modeling should be a part of routine application monitoring. Early, informative warnings on significant changes in application performance should help service providers to timely identify and prevent performance problems and their negative impact on the service. We propose a novel framework for automated anomaly detection and application change analysis. It is based on integration of two complementary techniques: i) a regression-based transaction model that reflects a resource consumption model of the application, and ii) an application performance signature that provides a compact model of run-time behavior of the application. The proposed integrated framework provides a simple and powerful solution for anomaly detection and analysis of essential performance changes in application behavior. An additional benefit of the proposed approach is its simplicity: it is not intrusive and is based on monitoring data that is typically available in enterprise production environments. Ludmila Cherkasova, Kivanc M. Ozonat, Ningfang Mi, Julie Symons, Evgenia Smirni |
DSN | 3 |
| 2008 | Scheduling for performance and availability in systems with temporal dependent workloadsabstractTemporal locality in workloads creates conditions in which a server, in order to remain available, should quickly process bursts of requests with large service requirements. In this paper, we show how to counteract the resulting peak congestions and maintain high availability by delaying selected requests that contribute to the temporal locality. We propose and evaluate SWAP, a measurement-based scheduling policy that approximates the shortest job first (SJF) scheduling without requiring any knowledge of job service times. We show that good service time estimates can be obtained from the temporal dependence structure of the workload and allow to closely approximate the behavior of SJF. Experimental results indicate that SWAP significantly improves system performability. In particular, we show that system capacity under SWAP is largely increased compared to first-come first-served (FCFS) scheduling and is highly-competitive with SJF, but without requiring a priori information of job service times. Ningfang Mi, Giuliano Casale, Evgenia Smirni |
DSN | 1 |
| 2008 | Enhancing data availability in disk drives through background activitiesabstractLatent sector errors in disk drives affect only a few data sectors. They occur silently and are detected only when the affected area is accessed again. If a latent error is detected while the storage system is operating under reduced redundancy, i.e., during a RAID rebuild, then data loss may occur. Various features such as scrubbing and intra-disk data redundancy are proposed to detect and/or recover from latent errors and avoid data loss. While such features enhance data availability in the storage system, their execution may cause performance degradation. In this paper, we evaluate the effectiveness of scrubbing and intra-disk data redundancy in improving data availability while the overall goal is to maintain user performance within predefined bounds. We show that by treating them as low priority background activities and scheduling them efficiently during idle times, these features remain performance-wise transparent to the storage system user while still improving data reliability. Detailed trace-driven simulations show that the mean time to data loss (MTTDL) improves by up to 5 orders of magnitude if these features are implemented independently. By scheduling concurrently both scrubbing and intra-disk parity updates during idle times in disk drives, MTTDL improves by as much as 8 orders of magnitude. Ningfang Mi, Alma Riska, Evgenia Smirni, Erik Riedel |
DSN | 1 |
| 2008 | Versatile models of systems using map queueing networksabstractAnalyzing the performance impact of temporal dependent workloads on hardware and software systems is a challenging task that yet must be addressed to enhance performance of real applications. For instance, existing matrix-analytic queueing models can capture temporal dependence only in systems that can be described by one or two queues, but the capacity planning of real multi-tier architectures requires larger models with arbitrary topology. To address the lack of a proper modeling technique for systems subject to temporal dependent workloads, we introduce a class of closed queueing networks where service times can have non-exponential distribution and accurately approximate temporal dependent features such as short or long range dependence. We describe these service processes using Markovian arrival processes (MAPs), which include the popular Markov-modulated Poisson processes (MMPPs) as special cases. Using a linear programming approach, we obtain for MAP closed networks tight upper and lower bounds for arbitrary performance indexes (e.g., throughput, response time, utilization). Numerical experiments indicate that our bounds achieve a mean accuracy error of 2% and promote our modeling approach for the accurate performance analysis of real multi-tier architectures. Giuliano Casale, Ningfang Mi, Evgenia Smirni |
IPDPS | 2 |
| 2008 | Burstiness in Multi-tier Applications: Symptoms, Causes, and New Models
Ningfang Mi, Giuliano Casale, Ludmila Cherkasova, Evgenia Smirni |
Middleware | 1 |
| 2008 | Analysis of application performance and its change via representative application signaturesabstractApplication servers are a core component of a multitier architecture that has become the industry standard for building scalable client-server applications. A client communicates with a service deployed as a multi-tier application via request-reply transactions. A typical server reply consists of the web page dynamically generated by the application server. The application server may issue multiple database calls while preparing the reply. Understanding the cascading effects of the various tasks that are sprung by a single request-reply transaction is a challenging task. Furthermore, significantly shortened time between new software releases further exacerbates the problem of thoroughly evaluating the performance of an updated application. We address the problem of efficiently diagnosing essential performance changes in application behavior in order to provide timely feedback to application designers and service providers. In this work, we propose a new approach based on an application signature that enables a quick performance comparison of the new application signature against the old one, while the application continues its execution in the production environment. The application signature is built based on new concepts that are introduced here, namely the transaction latency profiles and transaction signatures. These become instrumental for creating an application signature that accurately reflects important performance characteristics. We show that such an application signature is representative and stable under different workload characteristics. We also show that application signatures are robust as they effectively capture changes in transaction times that result from software updates. Application signatures provide a simple and powerful solution that can further be used for efficient capacity planning, anomaly detection, and provisioning of multi-tier applications in rapidly evolving IT environments. Ningfang Mi, Ludmila Cherkasova, Kivanc M. Ozonat, Julie Symons, Evgenia Smirni |
NOMS | 1 |
| 2008 | Bound analysis of closed queueing networks with workload burstinessabstractBurstiness and temporal dependence in service processes are often found in multi-tier architectures and storage devices and must be captured accurately in capacity planning models as these features are responsible of significant performance degradations. However, existing models and approximations for networks of first-come first-served (FCFS) queues with general independent (GI) service are unable to predict performance of systems with temporal dependence in workloads. Giuliano Casale, Ningfang Mi, Evgenia Smirni |
SIGMETRICS | 2 |
| 2008 | Performance-Guided Load (Un)balancing under Autocorrelated FlowsabstractSize-based policies have been shown in the literature to effectively balance the load and improve performance in cluster environments. Size-based policies assign jobs to servers based on the job size and their performance improvements are an outcome of separating ";short"; from ";long"; jobs, by avoiding having short jobs waiting behind long jobs for service. In this paper, we present evidence that performance improvements due to this separation quickly vanish if the arrival process to the cluster is autocorrelated. Based on our observations, we devise a new size-based policy called D_EQAL that still strives to separate jobs to servers according to job size but this separation is now biased by an effort to reduce performance loss due to autocorrelation in the arrival flows to each server. As a result of this bias, all servers may not be equally utilized (i.e., the load in the system may be ";unbalanced";), but performance benefits become significant. D_EQAL can be used on-line as it does not assume any a priori knowledge of the incoming workload. Extensive simulations show the effectiveness of D_EQAL under autocorrelated and uncorrelated arrival streams and illustrate that the policy successfully self- adjusts the degree of load unbalancing based on monitored performance measures. Qi Zhang 0012, Ningfang Mi, Alma Riska, Evgenia Smirni |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2007 | New Results on the Performance Effects of Autocorrelated Flows in SystemsabstractTemporal dependence within the workload of any computing or networking system has been widely recognized as a significant factor affecting performance. More specifically, burstiness, as a form of temporal dependency, is catastrophic for performance. We use the autocorrelation function in a workload flow to formalize burstiness and also to characterize temporal dependence within a flow. We present results from two application areas: load balancing in a homogeneous cluster environment and capacity planning in a multi-tiered e-commerce system. For the load balancing problem, we show that if autocorrelation exists in the arrival stream to the cluster, classic load balancing policies become ineffective and solutions that focus on "unbalancing" the load offer superior performance. For the case of multi-tiered systems, we show that if there is autocorrelation in the flows, we observe the surprising result that in spite of the fact that the bottleneck resource in the system is far from saturation and that the measured throughput and utilizations of other resources are also modest, user response times are very high. For multi-tired systems, this underutilization of resources falsely indicates that the system can sustain higher capacities. We present analysis of the above phenomena that aims at the development better scheduling policies under auto correlated flows. Evgenia Smirni, Qi Zhang 0012, Ningfang Mi, Alma Riska, Giuliano Casale |
IPDPS | 3 |
| 2007 | Efficient management of idleness in systemsabstractNo abstract available. Ningfang Mi, Alma Riska, Qi Zhang 0012, Evgenia Smirni, Erik Riedel |
SIGMETRICS | 1 |
| 2007 | Performance impacts of autocorrelated flows in multi-tiered systems
Ningfang Mi, Qi Zhang 0012, Alma Riska, Evgenia Smirni, Erik Riedel |
Perform. Evaluation | 1 |
| 2006 | Evaluating the Performability of Systems with Background JobsabstractAs most computer systems are expected to remain operational 24 hours a day, 7 days a week, they must complete maintenance work while in operation. This work is in addition to the regular tasks of the system and its purpose is to improve system reliability and availability. Nonetheless, additional work in the system, although labeled as best effort or low priority, still affects the performance of foreground tasks, especially if background/foreground work is non-preemptive. In this paper, we propose an analytic model to evaluate the performance trade-offs of the amount of background work that a storage system can sustain. The proposed model results in a quasi-birth-death (QBD) process that is analytically tractable. Detailed experimentation using a variety of workloads shows that under dependent arrivals both foreground and background performance strongly depends on system load. In contrast, if arrivals of foreground jobs are independent, performance sensitivity to load is reduced. The model identifies dependence in the arrivals of foreground jobs as an important characteristic that controls the decision of how much background load the system can accept to maintain high availability and performance gains Qi Zhang 0012, Ningfang Mi, Evgenia Smirni, Alma Riska, Erik Riedel |
DSN | 2 |
| 2006 | Load Unbalancing to Improve Performance under Autocorrelated TrafficabstractSize-based policies have been shown to successfully balance load and improve performance in homogeneous cluster environments where a dispatcher assigns a job to a server strictly based on the job size. While the success of size-based policies is based on separating jobs to different servers according to their sizes by avoiding the unfavorable performance effects of having short jobs been stuck behind long jobs, we show that their effectiveness quickly deteriorates in the presence of job arrivals that are characterized by correlation in their dependence structure. We propose a new policy that still strives to separate jobs according to their sizes, but this separation is biased by the effort to reduce the performance loss due to autocorrelation. As a result, not all servers are equally utilized (i.e., the load in the system becomes unbalanced) but the performance benefits of this load unbalancing are significant. The proposed policy can be used on-line, i.e., it does not assume any knowledge neither of the correlation structure of the arrival stream, nor of the job size distribution in the system. Via detailed trace-driven simulation we quantify the performance benefits of the proposed policy and we show that it can effectively self adjust its configuration parameters to improve performance under continuously changing workload conditions. Qi Zhang 0012, Ningfang Mi, Alma Riska, Evgenia Smirni |
ICDCS | 2 |
| 2006 | Farthest-point queries with geometric and combinatorial constraints
Ovidiu Daescu, Ningfang Mi, Chan-Su Shin, Alexander Wolff 0001 |
Comput. Geom. | 2 |
| 2005 | Polygonal path simplification with angle constraints
Danny Ziyi Chen, Ovidiu Daescu, John Hershberger 0001, Peter M. Kogge, Ningfang Mi, Jack Snoeyink |
Comput. Geom. | 5 |
| 2005 | Polygonal chain approximation: a query based approach
Ovidiu Daescu, Ningfang Mi |
Comput. Geom. | 2 |
| 2003 | Polygonal Path Approximation: A Query Based Approach
Ovidiu Daescu, Ningfang Mi |
ISAAC | 2 |