VLDB 2026 Research / reviewers in the wild / expert
Zhuozhao Li
dblp:148/1926
· DBLP profile ↗
52ranked-venue papers
19as first author
24since 2021 · last 2026
0000-0003-1903-6428ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 35 · 12 first-author · 18 since 2021Computer networks · 6 · 4 first-author · 3 since 2021Applied, interdisciplinary, general and emerging computing · 5 · 3 first-author · 1 since 2021Artificial intelligence and machine learning · 3 · 1 first-author · 1 since 2021Software engineering, systems software and programming languages · 2 · 2 since 2021Databases, data management, data science and information retrieval · 2 · 1 first-author
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | DynaGene: An Adaptive Scheduling Framework for AI-Coupled Workflows on Heterogeneous Resources
Xuxuan Yuan, Haifeng Ni, Ran Cheng 0004, Zhuozhao Li |
IPDPS | 4 |
| 2026 | Joint Privacy-Preserving Cloud Workload Prediction Based on Federated LearningabstractOrganizations around the world are increasingly using multiple clouds for critical workloads, which makes the accurate prediction of cloud workload a collaborative task across multiple cloud providers. Our experimental studies on four large-scale cloud workload datasets demonstrate that the workloads are highly heterogeneous, and different cloud providers have highly heterogeneous model preferences, which makes the above task challenging. Accordingly, we propose a joint privacy-preserving cloud workload prediction framework based on Federated Learning (FL), which enables the collaborative training of a workload prediction model across cloud providers with heterogeneous workloads and model preferences without exposing private workload data. To handle workload heterogeneity without privacy leakage, we first design a Generative Adversarial Network (GAN) based virtual workload dataset generation method that incorporates workloads’ temporal and resource-usage-plan-related features via temporal-aware feature calibration. Then, to handle model heterogeneity, we design a FL approach based on a hybrid local model composed of a homogeneous public model and a heterogeneous personal model. Finally, we design a Knowledge Distillation (KD) based approach to post-train the personal model with both private workload knowledge and shared workload knowledge learned from the trained global public model. Extensive experiments on various real-world workloads and actual implementations demonstrate that compared with the state-of-the-art, our framework improves workload prediction accuracy by 57.7% in average over all cloud providers. Li Yan 0004, Zhuozhao Li, Mingjun Kao, Huanbo Gao, Xingye Sun, Chao Shen 0001 |
IEEE Trans. Netw. | 3 |
| 2026 | P3Forecast: Personalized and Adaptive Cloud Workload Prediction via GAN-Based Federated Data AugmentationabstractTo comply with privacy regulations and self-protection from cloud outages, increasingly more users are leveraging multiple cloud providers for service deployment. However, due to the isolation of highly Non-Identically and Independently Distributed (Non-IID) workload datasets and deficiencies of existing methods as reflected in our experimental studies, no cloud providers can single-handedly capture the workload patterns of such users for accurate workload prediction. Accordingly, we proposeP3Forecast, aPersonalizedPrivacy-Preserving Cloud workload prediction framework based on Federated Generative Adversarial Networks (GANs), which allows cloud providers with Non-IID workload data to collaboratively train workload prediction models as preferred while protecting privacy. We first design a data synthesis quality assessment method based on Dynamic Time Warping (calledpattern-aware DTW), which is insusceptible to time series length and reliable for the comparison of temporal patterns. By usingpattern-aware DTWas the model aggregation weights, we adopt the Federated Learning (FL) of a GAN model for the augmentation of IID workload training datasets per cloud provider. Then, we further design a post-training method of local workload prediction models, which consists of a query mechanism based on comprehensive evaluation of data synthesis informativeness and an adaptive learning rate adjustment strategy for stable convergence. Extensive experiments driven by real-world workloads demonstrate that compared with the state-of-the-art,P3Forecastimproves workload prediction accuracy by 23.7%-64.5% on average over all cloud providers, while ensuring the fastest convergence in Federated GAN training. Xin Yong, Li Yan 0004, Yu Kuang, Zhuozhao Li, Chao Shen 0001, Xingwei Wang 0001 |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2025 | Wrath: Workload Resilience Across Task Hierarchies in Task-Based Parallel Programming FrameworksabstractFailures in Task-based Parallel Programming (TBPP) can severely degrade performance and result in incomplete or incorrect outcomes. Existing failure-handling approaches, including reactive, proactive, and resilient methods such as retry and checkpointing mechanisms, often apply uniform retry mechanisms regardless of the root cause of failures, failing to account for the unique characteristics of TBPP frameworks such as heterogeneous resource availability and task-level failures. To address these limitations, we propose Wrath, a novel systematic approach that categorizes failures based on the unique layered structure of TBPP frameworks and defines specific responses to address failures at different layers. Wrath combines a distributed monitoring system and a resilient module to collaboratively address different types of failures in real time. The monitoring system captures execution and resource information, reports failures, and profiles tasks across different layers of TBPP frameworks. The resilient module then categorizes failures and responds with appropriate actions, such as hierarchically retrying failed tasks on suitable resources. Evaluations demonstrate that Wrath significantly improves TBPP robustness, tripling the task success rate and maintaining an application success rate of over 90 % for resolvable failures. Additionally, Wrath can reduce the time to failure by$20 \%-50 \%$, allowing tasks that are destined to fail to be identified and fail more quickly. Zhuozhao Li, Valérie Hayot-Sasson, Haochen Pan, Maxime Gonthier, J. Gregory Pauloski, Ryan Chard, Kyle Chard, Ian T. Foster |
CCGrid | 2 |
| 2025 | DeFragS: Mitigating Resource Fragmentation in GPU Clusters Through Spatial-Temporal Scheduling
Haifeng Ni, Zhuozhao Li |
ICA3PP (2) | 3 |
| 2025 | ParaCOSM: A Parallel Framework for Continuous Subgraph MatchingabstractContinuous Subgraph Matching (CSM) has been widely studied, yet most single-threaded algorithms struggle with large query graphs. Existing CSM algorithms on CPU suffer from load imbalance in searching and sequential updates to the index structure. Haibin Lai, Site Fan, Zhuozhao Li |
ICPP | 4 |
| 2025 | Lias: Leveraging Performance Counters for Interference Quantification and Mitigation in Multi-processor SystemsabstractIn modern clouds, colocating latency-critical (LC) and best-effort (BE) applications on the same physical server can lead to performance interference, compromising the Quality of Service (QoS) for LC applications. While interference-aware schedulers often rely on runtime indicators to detect interference and adjust resources dynamically, existing approaches face two major limitations: i) performance counter collection itself may introduce non-negligible overhead, and ii) most techniques rely on a narrow set of metrics, limiting generality and offering poor visibility into root resource bottlenecks. To tackle these issues, we propose Lias, an interference-aware monitoring and scheduling system that dynamically collects performance counters to provide real-time insights into service performance. Lias introduces a novel composite metric, interference entropy, to capture resource contention across multiple performance counters, and leverages counter distributions to accurately localize resource pressure, enabling timely and precise resource reallocation with an average latency increase of less than 0.2 ms. Our experiments show that Lias reduces the average QoS violation ratio from 24.8% to 7.1%, while improving BE throughput by up to 28.4% compared to state-of-the-art approaches. Yangfan Qiao, Zhuozhao Li |
ICPP | 2 |
| 2025 | Heterogeneity-aware Task Scheduling based on Personalized Federated Reinforcement LearningabstractThe workload data generated in large-scale cloud environments is becoming increasingly complex, making collaborative training a promising approach for developing more efficient task schedulers. Considering privacy security and transfer costs, Federated Reinforcement Learning (FRL) emerges as a promising solution. However, our exploratory experiments demonstrate that the environmental heterogeneity contributes to performance degradation in FRL, which makes the above issue challenging. Accordingly, we propose a Personalized FRL method based on Dual-critic networks and Multi-head attention aggregator (PFRL-DM), which achieves the optimal scheduling policies by collaborative training on diverse workload data in heterogeneous environments without exposing private data. We initially introduce a novel Reinforcement Learning (RL) environment modeling, serving as a foundation for the collaborative training of the cloud scheduling agents. Then, we implement a dual-critic network Proximal Policy Optimization (PPO) algorithm for each client, effectively balancing the influence between global and local models on the agents. Furthermore, we integrate multi-head attention weights into the server-side aggregator to implement personalization. Extensive experiments on various real-world workloads have demonstrated that, compared to state-of-the-art algorithm MFPO, the proposed algorithm exhibits faster convergence, shorter response and completion times, and achieves the highest resource utilization. Additionally, the PFRL-DM algorithm constructs personalized models for each client, enabling greater adaptability in heterogeneous and hybrid workload environments. The codes for this paper can be found at https://github.com/liyan2015/PFRL-DM. Xin Yong, Li Yan 0004, Zhuozhao Li |
ICPP | 3 |
| 2025 | P3 Forecast: Personalized Privacy-Preserving Cloud Workload Prediction Based on Federated Generative Adversarial NetworksabstractTo comply with privacy regulations and selfprotection from cloud outages, increasingly more users are leveraging multiple cloud providers for service deployment. However, due to the isolation of highly Non-Identically and Independently Distributed (Non-IID) workload datasets and deficiencies of existing methods as reflected in our experimental studies, no cloud providers can single-handedly capture the workload patterns of such users for accurate workload prediction. Accordingly, we propose$\boldsymbol{P}^{3}$Forecast, a Personalized Privacy-Preserving Cloud workload prediction framework based on Federated Generative Adversarial Networks (GANs), which allows cloud providers with Non-IID workload data to collaboratively train workload prediction models as preferred while protecting privacy. We first design a data synthesis quality assessment method based on Dynamic Time Warping (called pattern-aware DTW), which is insusceptible to time series length and reliable for the comparison of temporal patterns. By using pattern-aware DTW as the model aggregation weights, we adopt the Federated Learning (FL) of a GAN model for the augmentation of IID workload training datasets per cloud provider. Then, we further design a post-training method of local workload prediction models, which consists of a query mechanism for extracting the most informative synthesized data for training dataset augmentation and a learning rate adjustment strategy for stable convergence. Extensive experiments driven by real-world workloads demonstrate that compared with the state-of-the-art,$P^{3}$Forecast improves workload prediction accuracy by$19.5 \%-46.7 \%$in average over all cloud providers, while ensuring the fastest convergence in Federated GAN training. The codes can be found at https://github.com/liyan2015/P3Forecast. Yu Kuang, Li Yan 0004, Zhuozhao Li |
IPDPS | 3 |
| 2025 | BOER: Enhancing Resource Utilization for Deep Learning Inference with Hybrid Spatial GPU SharingabstractMany inference systems leverage spatial multiplexing technologies, such as Multi-Process Service (MPS) and Multi-Instance GPU (MIG), to serve deep learning models concurrently on a single GPU. However, existing solutions suffer from interference under MPS and rigid partition sizes in MIG. To address these limitations, we propose BOER, a system that combines MPS atop MIG partitions to reduce interference and enhance GPU utilization. BOER identifies key challenges in integrating MPS with MIG and introduces a hierarchical scheduling framework that jointly determines model colocation, workload distribution, MIG partitioning, and MPS configurations, while minimizing resource fragmentation and MIG reconfiguration overhead. Since MPS interference is difficult to predict accurately, BOER avoids performance models and instead employs a Bayesian optimization with tailored acceleration strategies to efficiently explore the MPS configuration space. Evaluation on a real testbed demonstrates that BOER outperforms state-of-the-art spatial multiplexing solutions, improving inference throughput by up to 46.04%–77.19% while preserving Quality-of-Service requirements. Yuhang Wang 0012, Zhuozhao Li |
SC | 3 |
| 2025 | EvoX: A Distributed GPU-Accelerated Framework for Scalable Evolutionary ComputationabstractInspired by natural evolutionary processes, Evolutionary Computation (EC) has established itself as a cornerstone of Artificial Intelligence. Recently, with the surge in data-intensive applications and large-scale complex systems, the demand for scalable EC solutions has grown significantly. However, most existing EC infrastructures fall short of catering to the heightened demands of large-scale problem solving. While the advent of some pioneering GPU-accelerated EC libraries is a step forward, they also grapple with some limitations, particularly in terms of flexibility and architectural robustness. In response, we introduce EvoX: a computing framework tailored for automated, distributed, and heterogeneous execution of EC algorithms. At the core of EvoX lies a unique programming model to streamline the development of parallelizable EC algorithms, complemented by a computation model specifically optimized for distributed GPU acceleration. Building upon this foundation, we have crafted an extensive library comprising a wide spectrum of 50+ EC algorithms for both single-and multi-objective optimization. Furthermore, the library offers comprehensive support for a diverse set of benchmark problems, ranging from dozens of numerical test functions to hundreds of reinforcement learning tasks. Through extensive experiments across a range of problem scenarios and hardware configurations, EvoX demonstrates robust system and model performances. EvoX is open-source and accessible at: https://github.com/EMI-Group/EvoX. Beichen Huang, Ran Cheng 0004, Zhuozhao Li, Yaochu Jin, Kay Chen Tan |
IEEE Trans. Evol. Comput. | 3 |
| 2024 | MIGER: Integrating Multi-Instance GPU and Multi-Process Service for Deep Learning ClustersabstractModern NVIDIA GPUs, known for their powerful computational abilities, have been widely adopted by data centers. These GPUs often use space-sharing techniques, such as Multi-Process Service (MPS) and Multi-Instance GPU (MIG), to run multiple workloads on a GPU concurrently. However, our findings reveal that there are issues such as performance interference and inflexible resource size for these techniques when they are used individually. Zhuozhao Li |
ICPP | 3 |
| 2024 | UniFaaS: Programming across Distributed Cyberinfrastructure with Federated Function ServingabstractModern scientific applications are increasingly decomposable into individual functions that may be deployed across distributed and diverse cyberinfrastructure such as supercomputers, clouds, and accelerators. Such applications call for new approaches to programming, distributed execution, and function-level management. We present UniFaaS, a parallel programming framework that relies on a federated function-as-a-service (FaaS) model to enable composition of distributed, scalable, and high-performance scientific workflows, and to support fine-grained function-level management. UniFaaS provides a unified programming interface to compose dynamic task graphs with transparent wide-area data management. UniFaaS exploits an observe-predict-decide approach to efficiently map workflow tasks to target heterogeneous and dynamic resources. We propose a dynamic heterogeneity-aware scheduling algorithm that employs a delay mechanism and a re-scheduling mechanism to accommodate dynamic resource capacity. Our experiments show that UniFaaS can efficiently execute workflows across computing resources with minimal scheduling overhead. We show that UniFaaS can improve the performance of a real-world drug screening workflow by as much as 22.99% when employing an additional 19.48% of resources and a montage workflow by 54.41% when employing an additional 47.83% of resources across multiple distributed clusters, in contrast to using a single cluster. Ryan Chard, Yadu N. Babuji, Kyle Chard, Ian T. Foster, Zhuozhao Li |
IPDPS | 6 |
| 2024 | MAAD: A Distributed Anomaly Detection Architecture for Microservices SystemsabstractAnomaly detection plays a critical role in microservices systems by enabling system administrators to promptly detect and respond to anomalies. However, existing anomaly detection systems often necessitate the centralization of log and trace data from diverse system components and rely on resource-intensive statistical methods or deep learning models for analysis. This approach impedes real-time anomaly detection and requires a significant demand on computing resources. In this paper, we design a multi-agent-based, distributed anomaly detection architecture called MAAD to address the limitations. MAAD consists of a collection of agents that cooperate together to identify abnormal behaviors in a distributed manner. Each agent is deployed along with a single service and applies lightweight machine learning techniques to perform local anomaly detection based on its own logs, local context, and information extracted from its parent span service.To preserve the graph information in a microservices request, an agent can communicate essential features with each other, taking into account the collective patterns learned from the prior services. We evaluate the effectiveness of MAAD on two microservices datasets, TrainTicket and MicroSS, and show that MAAD achieved high precision (up to 95.8%) and recall (up to 99.6%), outperforming state-of-the-art centralized anomaly detection approaches. Compared to centralized approaches, MAAD reduces the amount of transferred data before anomaly detection by approximately 88%, facilitating real-time anomaly detection. Furthermore, the lightweight nature of MAAD allows for rapid anomaly detection with minimal impact on microservices systems. Compared to DeepLog, MAAD saves approximately 92% detection time without using GPU accelerators. Rongyuan Tan, Zhuozhao Li |
IPDPS | 2 |
| 2023 | gSampler: General and Efficient GPU-based Graph Sampling for Graph LearningabstractGraph sampling prepares training samples for graph learning and can dominate the training time. Due to the increasing algorithm diversity and complexity, existing sampling frameworks are insufficient in the generality of expression and the efficiency of execution. To close this gap, we conduct a comprehensive study on 15 popular graph sampling algorithms to motivate the design of gSampler, a general and efficient GPU-based graph sampling framework. gSampler models graph sampling using a general 4-step Extract-Compute-Select-Finalize (ECSF) programming model, proposes a set of matrix-centric APIs that allow to easily express complex graph sampling algorithms, and incorporates a data-flow intermediate representation (IR) that translates high-level API codes for efficient GPU execution. We demonstrate that implementing graph sampling algorithms with gSampler is easy and intuitive. We also conduct extensive experiments with 7 algorithms, 4 graph datasets, and 2 hardware configurations. The results show that gSampler introduces sampling speedups of 1.14--32.7× and an average speedup of 6.54×, compared to state-of-the-art GPU-based graph sampling systems such as DGL, which translates into an overall time reduction of over 40% for graph learning. gSampler is open-source at https://tinyurl.com/29twthd4. Ping Gong 0009, Renjie Liu 0001, Zunyao Mao, Zhenkun Cai, Xiao Yan 0002, Cheng Li 0001, Zhuozhao Li |
SOSP | 8 |
| 2022 | HiPS2022: The 2nd Workshop on High Performance Serverless ComputingabstractServerless computing presents an attractive model for general distributed computing as it focuses on abstracting the infrastructure required to execute an application. This workshop investigates the intersection between high performance computing and serverless computing, looking both at the high performance and distributed systems used to deliver serverless platforms and also at the use of serverless models for high performance and distributed systems. Yadu N. Babuji, Kyle Chard, Ian T. Foster, Zhuozhao Li |
HPDC | 4 |
| 2022 | Co-Scheduler: A Coflow-Aware Data-Parallel Job Scheduler in Hybrid Electrical/Optical Datacenter NetworksabstractTo support higher demand for datacenter networks, networking researchers have proposed hybrid electrical/optical datacenter networks (Hybrid-DCN) that leverages optical circuit switching (OCS) along with traditional electrical packet switching (EPS). However, due to the highreconfigurationdelay of OCS, OCS is used only for bulk data transfers between racks to amortize the reconfiguration delay. Existing job schedulers for data-parallel frameworks are not designed for Hybrid-DCN, since they neither place tasks to aggregate data traffic to take advantage of OCS, nor schedule tasks to minimize the Coflow completion time (CCT). In this paper, we describe the mismatch between existing job schedulers and the advanced Hybrid-DCN, introduce the requirements for the new scheduler, and present the implementation ofCo-scheduler, a job scheduler for data-parallel frameworks that aims to improve job performance by placing the tasks of jobs to aggregate enough data traffic to better leverage OCS to minimize the CCT in Hybrid-DCN. Specifically, for every job, Co-scheduler computesguidelineson how many racks to place the job’s input data and the job’s tasks. The guidelines are dynamically generated based on the real-time job characteristics or predictable job characteristics from prior runs, with the aim of leveraging OCS whenever possible and efficient and minimizing CCT of jobs. Co-scheduler then schedules the tasks of jobs based on the guidelines. We evaluate the effectiveness of Co-scheduler using trace-driven simulation. The evaluation demonstrates that Co-scheduler can improve makespan, average job completion time, and average CCT of a workload by up to 56%, 61%, and 79%, respectively, compared to the state-of-the-art schedulers. Zhuozhao Li, Haiying Shen |
IEEE/ACM Trans. Netw. | 1 |
| 2022 | A Study on the Impact of Memory DoS Attacks on Cloud Applications and Exploring Real-Time Detection SchemesabstractEven though memory denial-of-service attacks can cause severe performance degradations onco-locatedvirtual machines, a previous detection scheme against such attacks cannot accurately detect the attacks and also generates high detection delay and high performance overhead since it assumes that cache-related statistics of an application follow the same probability distribution at all times, which may not be true for all types of applications. In this paper, we present the experimental results showing the impacts of memory DoS attacks on different types of cloud-based applications. Based on these results, we propose two lightweight and responsive Statistical based Detection Schemes (SDS/B and SDS/P) that can detect such attacks accurately. SDS/B constructs a profile of normal range of cache-related statistics for all applications and use statistical methods to infer an attack when the real-time collected statistics exceed this normal range, while SDS/P exploits the increased periods of access patterns for periodic applications to infer an attack. Upon SDS, we further leverage deep neural network (DNN) techniques to design a DNN-based detection scheme that is general to various types of applications and more robust to adaptive attack scenarios. Our evaluation results show that SDS/B, SDS/P and DNN outperform the state-of-the-art detection scheme, e.g., with 65% higher specificity, 40% shorter detection delay, and 7% less performance overhead. We also discuss how to use SDS and DNN-based detection schemes under different situations. Zhuozhao Li, Tanmoy Sen, Haiying Shen, Mooi Choo Chuah |
IEEE/ACM Trans. Netw. | 1 |
| 2022 | $f$funcX: Federated Function as a Service for ScienceabstractƒuncX is a distributed function as a service (FaaS) platform that enables flexible, scalable, and high performance remote function execution. Unlike centralized FaaS systems, ƒuncX decouples the cloud-hosted management functionality from the edge-hosted execution functionality. ƒuncX's endpoint software can be deployed, by users or administrators, on arbitrary laptops, clouds, clusters, and supercomputers, in effect turning them into function serving systems. ƒuncX's cloud-hosted service provides a single location for registering, sharing, and managing both functions and endpoints. It allows for transparent, secure, and reliable function execution across the federated ecosystem of endpoints—enabling users to route functions to endpoints based on specific needs. ƒuncX uses containers (e.g., Docker, Singularity, and Shifter) to provide common execution environments across endpoints. ƒuncX implements various container management strategies to execute functions with high performance and efficiency on diverse ƒuncX endpoints. ƒuncX also integrates with an in-memory data store and Globus for managing data that may span endpoints. We motivate the need for ƒuncX, present our prototype design and implementation, and demonstrate, via experiments on two supercomputers, that ƒuncX can scale to more than 130000 concurrent workers. We show that ƒuncX's container warming-aware routing algorithm can reduce the completion time for 3,000 functions by up to 61% compared to a randomized algorithm and the in-memory data store can speed up data transfers by up to 3x compared to a shared file system. Zhuozhao Li, Ryan Chard, Yadu N. Babuji, Ben Galewsky, Tyler J. Skluzacek, Kirill Nagaitsev, Anna Woodard, Ben Blaiszik, Josh Bryan, Daniel S. Katz, Ian T. Foster, Kyle Chard |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2021 | Federated Function as a Service for eScienceabstractThe function as a service paradigm aims to abstract the complexities of managing computing infrastructure for users. While adoption in industry has been swift, we have yet to see widespread adoption in academia. This is in part due to barriers such as the need to access large research data, diverse hardware requirements, monolithic code bases, and existing systems available to researchers. We describe funcX, a federated functionas-a-service platform that addresses important requirements for use of FaaS in research computing. We outline how funcX has been used in early science deployments. Yadu N. Babuji, Josh Bryan, Ryan Chard, Kyle Chard, Ian T. Foster, Ben Galewsky, Daniel S. Katz, Zhuozhao Li |
e-Science | 8 |
| 2021 | A Serverless Framework for Distributed Bulk Metadata ExtractionabstractWe introduce Xtract, an automated and scalable system for bulk metadata extraction from large, distributed research data repositories. Xtract orchestrates the application of metadata extractors to groups of files, determining which extractors to apply to each file and, for each extractor and file, where to execute. A hybrid computing model, built on the funcX federated FaaS platform, enables Xtract to balance tradeoffs between extraction time and data transfer costs by dispatching each extraction task to the most appropriate location. Experiments on a range of clouds and supercomputers show that Xtract can efficiently process multi-million-file repositories by orchestrating the concurrent execution of container-based extractors on thousands of nodes. We highlight the flexibility of Xtract by applying it to a large, semi-curated scientific data repository and to an uncurated scientific Google Drive repository. We show that by remotely orchestrating metadata extraction across decentralized storage and compute nodes, Xtract can process large repositories in 50% of the time it takes just to transfer the same data to a machine within the same computing facility. We also show that when transferring data is necessary (e.g., no local compute is available), Xtract can scale to process files as fast as they are received, even over a multi-GB/s network. Tyler J. Skluzacek, Ryan Wong 0002, Zhuozhao Li, Ryan Chard, Kyle Chard, Ian T. Foster |
HPDC | 3 |
| 2021 | IMPECCABLE: Integrated Modeling PipelinE for COVID Cure by Assessing Better LEadsabstractThe drug discovery process currently employed in the pharmaceutical industry typically requires about 10 years and $2–3 billion to deliver one new drug. This is both too expensive and too slow, especially in emergencies like the COVID-19 pandemic. In silico methodologies need to be improved both to select better lead compounds, so as to improve the efficiency of later stages in the drug discovery protocol, and to identify those lead compounds more quickly. No known methodological approach can deliver this combination of higher quality and speed. Here, we describe an Integrated Modeling PipEline for COVID Cure by Assessing Better LEads (IMPECCABLE) that employs multiple methodological innovations to overcome this fundamental limitation. We also describe the computational framework that we have developed to support these innovations at scale, and characterize the performance of this framework in terms of throughput, peak performance, and scientific results. We show that individual workflow components deliver 100 × to 1000 × improvement over traditional methods, and that the integration of methods, supported by scalable infrastructure, speeds up drug discovery by orders of magnitudes. IMPECCABLE has screened ∼ 1011 ligands and has been used to discover a promising drug candidate. These capabilities have been used by the US DOE National Virtual Biotechnology Laboratory and the EU Centre of Excellence in Computational Biomedicine. Aymen Alsaadi, Dario Alfè, Yadu N. Babuji, Agastya Bhati, Ben Blaiszik, Alex Brace, Thomas S. Brettin, Kyle Chard, Ryan Chard, Austin Clyde, Peter V. Coveney, Ian T. Foster, Tom Gibbs, Shantenu Jha, Kristopher Keipert, Dieter Kranzlmüller, Thorsten Kurth, Hyungro Lee, Zhuozhao Li, Gerald Mathias, André Merzky, Alexander Partin, Arvind Ramanathan, Ashka Shah, Abraham C. Stern, Rick L. Stevens, Mikhail Titov, Anda Trifan, Aristeidis Tsaris, Matteo Turilli, Huub J. J. Van Dam, Shunzhou Wan, David Wifling, Junqi Yin |
ICPP | 19 |
| 2021 | Lightweight Function Monitors for Fine-Grained Management in Large Scale Python ApplicationsabstractPython has become a widely used programming language for research, not only for small one-off analyses, but also for complex application pipelines running at supercomputer-scale. Modern parallel programming frameworks for Python present users with a more granular unit of management than traditional Unix processes and batch submissions: the Python function. We review the challenges involved in running native Python functions at scale, and present techniques for dynamically determining a minimal set of dependencies and for assembling a lightweight function monitor (LFM) that captures the software environment and manages resources at the granularity of single functions. We evaluate these techniques in a range of environments, from campus cluster to supercomputer, and show that our advanced dependency management planning and dynamic resource management methods provide superior performance and utilization relative to coarser-grained management approaches, achieving several-fold decrease in execution time for several large Python applications. Timothy Shaffer, Zhuozhao Li, Benjamín Tovar, Yadu N. Babuji, T. J. Dasso, Zoe Surma, Kyle Chard, Ian T. Foster, Douglas Thain |
IPDPS | 2 |
| 2021 | DLHub: Simplifying publication, discovery, and use of machine learning models in science
Zhuozhao Li, Ryan Chard, Logan T. Ward, Kyle Chard, Tyler J. Skluzacek, Yadu N. Babuji, Anna Woodard, Steven Tuecke, Ben Blaiszik, Michael J. Franklin, Ian T. Foster |
J. Parallel Distributed Comput. | 1 |
| 2020 | Time-Efficient Geo-Obfuscation to Protect Worker Location Privacy over Road Networks in Spatial CrowdsourcingabstractTo promote cost-effective task assignment in Spatial Crowdsourcing (SC), workers are required to report their location to servers, which raises serious privacy concerns. As a solution, geo-obfuscation has been widely used to protect the location privacy of SC workers, where workers are allowed to report perturbed location instead of the true location. Yet, most existing geo-obfuscation methods consider workers? mobility on a 2 dimensional (2D) plane, wherein workers can move in arbitrary directions. Unfortunately, 2D-based geo-obfuscation is likely to generate high traveling cost for task assignment over roads, as it cannot accurately estimate the traveling costs distortion caused by location obfuscation. In this paper, we tackle the SC worker location privacy problem over road networks. Considering the network-constrained mobility features of workers, we describe workers? mobility by a weighted directed graph, which considers the dynamic traffic condition and road network topology. Based on the graph model, we design a geo-obfuscation (GO) function for workers to maximize the workers? overall location privacy without compromising the task assignment efficiency. We formulate the problem of deriving the optimal GO function as a linear programming (LP) problem. By using the angular block structure of the LP's constraint matrix, we apply Dantzig-Wolfe decomposition to improve the time-efficiency of the GO function generation. Our experimental results in the real-trace driven simulation and the real-world experiment demonstrate the effectiveness of our approach in terms of both privacy and task assignment efficiency. Chenxi Qiu, Anna Cinzia Squicciarini, Zhuozhao Li, Ce Pang, Li Yan 0004 |
CIKM | 3 |
| 2020 | funcX: A Federated Function Serving Fabric for ScienceabstractExploding data volumes and velocities, new computational methods and platforms, and ubiquitous connectivity demand new approaches to computation in the sciences. These new approaches must enable computation to be mobile, so that, for example, it can occur near data, be triggered by events (e.g., arrival of new data), be offloaded to specialized accelerators, or run remotely where resources are available. They also require new design approaches in which monolithic applications can be decomposed into smaller components, that may in turn be executed separately and on the most suitable resources. To address these needs we present funcX---a distributed function as a service (FaaS) platform that enables flexible, scalable, and high performance remote function execution. funcX's endpoint software can transform existing clouds, clusters, and supercomputers into function serving systems, while funcX's cloud-hosted service provides transparent, secure, and reliable function execution across a federated ecosystem of endpoints. We motivate the need for funcX with several scientific case studies, present our prototype design and implementation, show optimizations that deliver throughput in excess of 1 million functions per second, and demonstrate, via experiments on two supercomputers, that funcX can scale to more than more than 130 000 concurrent workers. Ryan Chard, Yadu N. Babuji, Zhuozhao Li, Tyler J. Skluzacek, Anna Woodard, Ben Blaiszik, Ian T. Foster, Kyle Chard |
HPDC | 3 |
| 2020 | Impact of Memory DoS Attacks on Cloud Applications and Real-Time Detection SchemesabstractIn this poster, we present measurement studies of the impact of memory DoS attacks on different types of cloud-based applications. Based on the observations, we propose a lightweight, responsive Statistical based Detection Scheme (SDS) that can detect such attacks accurately. Our initial evaluation results show that SDS outperforms the state-of-the-art detection scheme up to 2% higher recall, up to 65% higher specificity, and up to 40% shorter detection delay. Zhuozhao Li, Tanmoy Sen, Haiying Shen, Mooi Choo Chuah |
ICDCS | 1 |
| 2020 | Impact of Memory DoS Attacks on Cloud Applications and Real-Time Detection SchemesabstractEven though memory-based denial-of-service attacks can cause severe performance degradations on co-located virtual machines, a previous detection scheme against such attacks cannot accurately detect the attacks and also generates high detection delay and high performance overhead since it assumes that cache-related statistics of an application follow the same probability distribution at all times, which may not be true for all types of applications. In this paper, we present the experimental results showing the impacts of memory DoS attacks on different types of cloud-based applications. Based on these results, we propose two lightweight, responsive Statistical based Detection Schemes (SDS/B and SDS/P) that can detect such attacks accurately. SDS/B constructs a profile of normal range of cache-related statistics for all applications and use statistical methods to infer an attack when the real-time collected statistics exceed this normal range, while SDS/P exploits the increased periods of access patterns for periodic applications to infer an attack. Our evaluation results show that SDS/B and SDS/P outperform the state-of-the-art detection scheme, e.g., with 65% higher specificity, 40% shorter detection delay, and 7% less performance overhead. Zhuozhao Li, Tanmoy Sen, Haiying Shen, Mooi Choo Chuah |
ICPP | 1 |
| 2019 | ParaOpt: Automated Application Parameterization and Optimization for the CloudabstractThe variety of instance types available on cloud platforms offers enormous flexibility to match the requirements of applications with available resources. However, selecting the most suitable instance type and configuring an application to optimally execute on that instance type can be complicated and time-consuming. For example, application parallelism flags must match available cores and problem sizes must be tuned to match available memory. As the search space of application configurations can be enormous, we propose an automated approach, called ParaOpt, to automatically explore and tune application configurations on arbitrary cloud instances. ParaOpt supports arbitrary applications, enables use of custom optimization methods, and can be configured with different optimization targets such as runtime and cost. We evaluate ParaOpt by optimizing genomics, molecular dynamics, and machine learning applications with four types of optimizers. We show with as few as 15 parameterized executions of an application, representing between 1.2%-26.7% of the search space, that ParaOpt is able to identify the optimal configuration in 32.7% of experiments and a near-optimal configuration in 83.2% of cases. As a result of using near-optimal configurations, ParaOpt reduces overall execution time by up to 85.8% when compared with using the default configuration. Ian T. Foster, Ted Summer, Zhuozhao Li, Anna Woodard, Ryan Chard, Matt Baughman, Yadu N. Babuji, Kyle Chard, Jason Pitt |
CloudCom | 4 |
| 2019 | Parsl: Pervasive Parallel Programming in PythonabstractHigh-level programming languages such as Python are increasingly used to provide intuitive interfaces to libraries written in lower-level languages and for assembling applications from various components. This migration towards orchestration rather than implementation, coupled with the growing need for parallel computing (e.g., due to big data and the end of Moore's law), necessitates rethinking how parallelism is expressed in programs. Here, we present Parsl, a parallel scripting library that augments Python with simple, scalable, and flexible constructs for encoding parallelism. These constructs allow Parsl to construct a dynamic dependency graph of components that it can then execute efficiently on one or many processors. Parsl is designed for scalability, with an extensible set of executors tailored to different use cases, such as low-latency, high-throughput, or extreme-scale execution. We show, via experiments on the Blue Waters supercomputer, that Parsl executors can allow Python scripts to execute components with as little as 5 ms of overhead, scale to more than 250000 workers across more than 8000 nodes, and process upward of 1200 tasks per second. Other Parsl features simplify the construction and execution of composite programs by supporting elastic provisioning and scaling of infrastructure, fault-tolerant execution, and integrated wide-area data management. We show that these capabilities satisfy the needs of many-task, interactive, online, and machine learning applications in fields such as biology, cosmology, and materials science. Yadu N. Babuji, Anna Woodard, Zhuozhao Li, Daniel S. Katz, Ben Clifford, Lukasz Lacinski, Ryan Chard, Justin M. Wozniak, Ian T. Foster, Michael Wilde, Kyle Chard |
HPDC | 3 |
| 2019 | Accelerating Big Data Analytics Using Scale-Up/Out Heterogeneous ClustersabstractProduction data analytic workloads typically consist of a majority of jobs with small input data sizes and a small number of jobs with large input data sizes. Recent works advocate scale-up/scale-out heterogeneous clusters (in short Hybrid clusters) to handle these heterogeneous workloads, since scaleup machines (i.e., adding more resources to a single machine) can process small jobs faster than simply scaling out the cluster with cheap machines. However, there are several challenges for job placement and data placement to implement such a Hybrid cluster. In this paper, we propose a job placement strategy and a data placement strategy to solve the challenges. The job placement strategy places a job to either scale-up or scale-out machines based on the job's characteristics, and migrates jobs from scale-up machines to under-utilized scale-out machines to achieve load balance. The data placement strategy allocates data replicas in the two types of machines accordingly to increase the data locality in Hybrid cluster. We implemented a Hybrid cluster on Apache YARN, and evaluated its performance using a Facebook production workload. With our proposed strategies, a Hybrid cluster can reduce the makespan of the workload up to 37% and the median job completion time up to 60%, compared to traditional scale-out clusters with state-of-the-art schedulers. Zhuozhao Li, Haiying Shen, Lee Ward |
ICCCN | 1 |
| 2019 | Road Gradient Estimation Using Smartphones: Towards Accurate Estimation on Fuel Consumption and Air Pollution Emission on RoadsabstractAccurate estimations on vehicle fuel consumption and pollution emission on roads are important for vehicle velocity optimization and driving route planning. Existing methods for such estimations only consider vehicle driving speed and acceleration but neglect the influence of road gradient. This is mainly because the road gradients for most road networks are not available and none of existing methods for road gradient estimation can be conducted inexpensively in practice and keep high road gradient estimation accuracy simultaneously. Thus, how to estimate the road gradient conveniently and accurately is an important but challenging problem. To handle this challenge, we propose a new road gradient estimation system which estimates the road gradient only using a smartphone. When a vehicle is driving, a smartphone in the vehicle continuously measures vehicle states (velocity, acceleration, steering rate, position), which are used to estimate the road gradient. To eliminate measuring noise and drift noise, the deviation between the measured value and estimated value is used to adjust the estimated value. Since measured vehicle states when a vehicle changes lane adversely influence the accuracy of road gradient estimation, we design lane change detection to eliminate such influences. Finally, given a group of road gradient estimates for a given route, we use the track fusion algorithm to further eliminate measuring noise and drift noise and improve road gradient estimation accuracy. We conducted driving experiments in a city area to evaluate our system. The experimental results show that our system's estimation error is reduced by 22% compared with existing methods. The results also demonstrate the accuracy of our lane change detection. Finally, we integrated the road gradient values into vehicle fuel consumption and air pollution emission model to estimate fuel consumption and air pollution emission and found that the estimation values increase by 33.4% compared with the values without considering road gradient. Liuwang Kang, Haiying Shen, Zhuozhao Li |
ICDCS | 3 |
| 2019 | Co-scheduler: Accelerating Data-Parallel Jobs in Datacenter Networks with Optical Circuit SwitchingabstractThe optical circuit switch (OCS) in recently proposed hybrid electrical/optical datacenter networks (Hybrid-DCN) can only be used to transfer large flows (i.e., flows with a large size of data). Current job schedulers for data-parallel frameworks are not suitable for Hybrid-DCN, since they neither place tasks to aggregate data traffic to take advantage of OCS nor schedule tasks to minimize the Coflow completion time (CCT). In this paper, we propose Co-scheduler, a job scheduler for dataparallel frameworks that aims to improve job performance by attempting to place the tasks of a job to aggregate enough data traffic to take advantage of OCS and minimize the CCT in Hybrid-DCN. Specifically, for each job, Co-scheduler computes a guideline on the number of racks to place the job's input data and to run the job's map tasks, so that the job can potentially take full advantage of OCS to transfer its data. When the map tasks of a job complete, Co-scheduler computes all the possible schedules of the reduce tasks of the job. Each possible schedule includes the number of racks to schedule the reduce tasks that enables the job to use OCS to transfer its data, and the number of reduce tasks to place on each of the racks that minimizes CCT of the job. Next, Co-scheduler selects a best schedule among all the possible schedules so that the job completion time is minimized. Finally, Co-scheduler schedules the map tasks and reduce tasks of the job based on the computed guideline and best schedule. The evaluation demonstrates that compared to the state-of-theart schedulers, Co-scheduler achieves performance improvements on makespan, average job completion time, and average CCT by up to 51.2%, 54.6% and 73.6%, respectively. Zhuozhao Li, Haiying Shen |
ICDCS | 1 |
| 2019 | JobPacker: Job Scheduling for Data-Parallel Frameworks with Hybrid Electrical/Optical Datacenter NetworksabstractIn spite of many advantages of hybrid electrical/optical datacenter networks (Hybrid-DCN), current job schedulers for data-parallel frameworks are not suitable for Hybrid-DCN, since the schedulers do not aggregate data traffic to facilitate using optical circuit switch (OCS). In this paper, we propose JobPacker, a job scheduler for data-parallel frameworks in Hybrid-DCN that aims to take full advantage of OCS to improve job performance. JobPacker aggregates the data transfers of a job in order to use OCS to improve data transfer efficiency. It first explores the tradeoff between parallelism and traffic aggregation for each shuffle-heavy recurring job, and then generates an offline schedule including which racks to run each job and the sequence to run the recurring jobs in each rack that yields the best performance. It has a new sorting method to prioritize recurring jobs in offline-scheduling to prevent high resource contention while fully utilizing cluster resources. In real-time scheduler, JobPacker uses the offline schedule to guide the data placement and schedule recurring jobs, and schedules non-recurring jobs to the idle resources not assigned to recurring jobs. Trace-driven simulation and GENI-based emulation show that JobPacker reduces the makespan up to 49% and the median completion time up to 43%, compared to the state-of-the-art schedulers in Hybrid-DCN. Zhuozhao Li, Haiying Shen |
ICPP | 1 |
| 2019 | DLHub: Model and Data Serving for ScienceabstractWhile the Machine Learning (ML) landscape is evolving rapidly, there has been a relative lag in the development of the “learning systems” needed to enable broad adoption. Furthermore, few such systems are designed to support the specialized requirements of scientific ML. Here we present the Data and Learning Hub for science (DLHub), a multi-tenant system that provides both model repository and serving capabilities with a focus on science applications. DLHub addresses two significant shortcomings in current systems. First, its self-service model repository allows users to share, publish, verify, reproduce, and reuse models, and addresses concerns related to model reproducibility by packaging and distributing models and all constituent components. Second, it implements scalable and low-latency serving capabilities that can leverage parallel and distributed computing resources to democratize access to published models through a simple web interface. Unlike other model serving frameworks, DLHub can store and serve any Python 3-compatible model or processing function, plus multiple-function pipelines. We show that relative to other model serving systems including TensorFlow Serving, SageMaker, and Clipper, DLHub provides greater capabilities, comparable performance without memoization and batching, and significantly better performance when the latter two techniques can be employed. We also describe early uses of DLHub for scientific applications. Ryan Chard, Zhuozhao Li, Kyle Chard, Logan T. Ward, Yadu N. Babuji, Anna Woodard, Steven Tuecke, Ben Blaiszik, Michael J. Franklin, Ian T. Foster |
IPDPS | 2 |
| 2018 | A Network-Aware Scheduler in Data-Parallel Clusters for High PerformanceabstractIn spite of many shuffle-heavy jobs in current commercial data-parallel clusters, few previous studies have considered the network traffic in the shuffle phase, which contains a large amount of data transfers and may adversely affect the cluster performance. In this paper, we propose a network-aware scheduler (NAS) that handles two main challenges associated with the shuffle phase for high performance: i) balancing cross-node network load, and ii) avoiding and reducing cross-rack network congestion. NAS consists of three main mechanisms: i) map task scheduling (MTS), ii) congestion-avoidance reduce task scheduling (CA-RTS) and iii) congestion-reduction reduce task scheduling (CR-RTS). MTS constrains the shuffle data on each node when scheduling the map tasks to balance the cross-node network load. CA-RTS distributes the reduce tasks for each job based on the distribution of its shuffle data among the racks in order to minimize cross-rack traffic. When the network is congested, CR-RTS schedules reduce tasks that generate negligible shuffle traffic to reduce the congestion. We implemented NAS in Hadoop on a cluster. Our trace-driven simulation and real cluster experiment demonstrate the superior performance of NAS on improving the throughput (up to 62%), reducing the average job execution time (up to 44%) and reducing the cross-rack traffic (up to 40%) compared with state-of-the-art schedulers. Zhuozhao Li, Haiying Shen, Ankur Sarker |
CCGrid | 1 |
| 2018 | PageRankVM: A PageRank Based Algorithm with Anti-Collocation Constraints for Virtual Machine Placement in Cloud DatacentersabstractThere is a dramatic increase in the variety of virtual machines (VMs) and complexity of VM placement problems in clouds. Previous VM placement approaches attempt to accommodate more VMs efficiently on fewer PMs by balancing the resource usages across multiple dimensions. However, these approaches are not sufficiently accurate in measuring the quality of the PMs in terms of fully utilizing PM resource and having the potential to accommodate more VMs. Therefore, it is critical to design a new method that can more accurately measure the probability of a PM of fully utilizing its resources after accommodating a given VM with the consideration of different types of VMs. In addition, anti-collocation constraints must be handled efficiently. We propose a PageRank based VM placement algorithm with anti-collocation constraints (PageRankVM). PageRankVM defines the best PM resource usage profile, which means that the PM has full resource utilization for every resource dimension, and then ranks PM profiles according to their convergence of transferring (by accommodating VMs) to the best profile. PageRankVM then places a given VM to the PM based on the ranks of the resulted PM profiles after accommodating the VM with the consideration of anti-collocation constraints. Compared to previous approaches, PageRankVM effectively measures the ability of different PM profiles to reach the best profiles by accommodating a given VM, and hence differentiates the effectiveness of different VM placement decisions. We conducted extensive trace-driven simulation and GENI testbed experiments and demonstrated that PageRankVM has superior performance compared with other methods in terms of reducing the number of PMs, the energy consumption, the number of VM migrations, and the service level objective (SLO) violations. Zhuozhao Li, Haiying Shen, Cole Miles |
ICDCS | 1 |
| 2018 | Approaches for Resilience against Cascading Failures in Cloud DatacentersabstractIn a modern cloud datacenter, a cascading failure will cause many Service Level Objective (SLO) violations. In a cascading failure, when a set of physical machines (PMs) in a failure domain are failed, their workloads are transferred to the PMs in another failure domain to continue. However, the new domain receiving additional workloads may become overloaded due to the resource oversubscription feature in the cloud, which easily leads to domain failures and subsequent workload transfer to other domains. This process repeats and a cascading failure is created finally. However, few previous methods can effectively handle the cascading failures. To handle this problem, we propose a Cascading Failure Resilience System (CFRS), which incorporates three methods: Overload-Avoidance VM Reassignment (OAVR), VM backup set placement (VMset) and Dynamic Oversubscription Ratio Adjustment (DOA). The experiments in trace-driven simulation show that CFRS outperforms other comparison methods in terms of the number of domain failures, the number of failed PMs and the number of SLO violations. Haoyu Wang 0003, Haiying Shen, Zhuozhao Li |
ICDCS | 3 |
| 2018 | Analysis of Knowledge Sharing Activities on a Social Network Incorporated Discussion Forum: A Case Study of DISboardsabstractDISboards is a discussion forum that provides a platform for knowledge sharing on planning and resources for Disney-related travel (Disney World, Disney Cruise Line, etc.). Since no previous work has been devoted to studying the online social networks (SNs) in the forums, we examine the SN and knowledge sharing activities in DISboards as a case study of discussion forums. Based on a large amount of data collected, we provide an in-depth study of DISboards. In particular, we analyzed SN structure, effect of SN in the forum, category characteristics and so on. We found that users with more friends are generally more active in the forum; teens are more active and constitute a significant part of the SN. We clustered the selected categories (e.g., resorts, dining, and hotels) into three groups: report, fact, discussion, and characterized their properties. Most users focus narrowly on only a few categories, while very few users participate in many categories. The development of SN should be able to attract more users to involve in the forum. We believe that the results presented in this paper are crucial in understanding SN and knowledge sharing in the forums. The paper also gives an instruction for the enhancement of SNs to incentivize users' activeness in the forums. Zhuozhao Li, Harrison Chandler, Haiying Shen |
IEEE Trans. Big Data | 1 |
| 2017 | Job scheduling for data-parallel frameworks with hybrid electrical/optical datacenter networksabstractIn spite of many advantages of hybrid electrical/optical datacenter networks (Hybrid-DCN), current job schedulers for data-parallel frameworks are not suitable for Hybrid-DCN, since the schedulers do not aggregate data traffic to facilitate using optical circuit switch (OCS). We propose SchedOCS, a job scheduler for data-parallel frameworks in Hybrid-DCN that aims to take full advantage of the OCS to improve the job performance. Zhuozhao Li, Haiying Shen |
SoCC | 1 |
| 2017 | Opportunistic Energy Sharing Between Power Grid and Electric Vehicles: A Game Theory-Based Pricing PolicyabstractElectric vehicles (EVs) have great potential to reduce dependency on fossil fuels. The recent surge in the development of online EV (OLEV) will help to address the drawbacks associated with current generation EVs, such as the heavy and expensive batteries. OLEVs are integrated with the smart grid of power infrastructure through a wireless power transfer system (WPT) to increase the driving range of the OLEV. However, the integration of OLEVs with the grid creates a tremendous load for the smart grid. The demand of a power grid changes over time and the price of power is not fixed throughout the day. There should be some congestion avoidance and load balancing policy implications to ensure quality of services for OLEVs. In this paper, first, we conduct an analysis to show the existence of unpredictable power load and congestion because of OLEVs. We use the Simulation for Urban MObility tool and hourly traffic counts of a road section of the New York City to analyze the amount of energy OLEVs can receive at different times of the day. Then, we present a game theory based on a distributed power schedule framework to find the optimal schedule between OLEVs and smart grid. In the proposed framework, OLEVs receive the amount of power charging from the smart grid based on a power payment function which is updated using best response strategy. We prove that the updated power requests converge to the optimal power schedule. In this way, the smart grid maximizes the social welfare of OLEVs, which is defined as mixed consideration of total satisfaction and its power charging cost. Finally, we verify the performance of our proposed pricing policy under different scenarios in a simulation study. Ankur Sarker, Zhuozhao Li, William Kolodzey, Haiying Shen |
ICDCS | 2 |
| 2017 | Measuring Scale-Up and Scale-Out Hadoop with Remote and Local File Systems and Selecting the Best PlatformabstractMapReduce is a popular computing model for parallel data processing on large-scale datasets, which can vary from gigabytes to terabytes and petabytes. Though Hadoop MapReduce normally uses Hadoop Distributed File System (HDFS) local file system, it can be configured to use a remote file system. Then, an interesting question is raised: for a given application, which is the best running platform among the different combinations of scale-up and scale-out Hadoop with remote and local file systems. However, there has been no previous research on how different types of applications (e.g., CPU-intensive, data-intensive) with different characteristics (e.g., input data size) can benefit from the different platforms. Thus, in this paper, we conduct a comprehensive performance measurement of different applications on scale-up and scale-out clusters configured with HDFS and a remote file system (i.e., OFS), respectively. We identify and study how different job characteristics (e.g., input data size, the number of file reads/writes, and the amount of computations) affect the performance of different applications on the different platforms. Based on the measurement results, we also propose a performance prediction model to help users select the best platforms that lead to the minimum latency. Our evaluation using a Facebook workload trace demonstrates the effectiveness of our prediction model. This study is expected to provide a guidance for users to choose the best platform to run different applications with different characteristics in the environment that provides both remote and local storage, such as HPC cluster and cloud environment. Zhuozhao Li, Haiying Shen |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2017 | An Exploration of Designing a Hybrid Scale-Up/Out Hadoop Architecture Based on Performance MeasurementsabstractScale-up machines perform better for jobs with small and median (KB, MB) data sizes, while scale-out machines perform better for jobs with large (GB, TB) data size. Since a workload usually consists of jobs with different data size levels, we propose building a hybrid Hadoop architecture that includes both scale-up and scale-out machines, which however is not trivial. The first challenge is workload data storage. Thousands of small data size jobs in a workload may overload the limited local disks of scale-up machines. Jobs from scale-up and scale-out machines may both request the same set of data, which leads to data transmission between the machines. The second challenge is to automatically schedule jobs to either scale-up or scale-out cluster to achieve the best performance. We conduct a thorough performance measurement of different applications on scale-up and scale-out clusters, configured with Hadoop Distributed File System (HDFS) and a remote file system (i.e., OFS), respectively. We find that using OFS rather than HDFS can solve the data storage challenge. Also, we identify the factors that determine the performance differences on the scale-up and scale-out clusters and their cross points to make the choice. Accordingly, we design and implement the hybrid scale-up/out Hadoop architecture. Our trace-driven experimental results show that our hybrid architecture outperforms both the traditional Hadoop architecture with HDFS and with OFS in terms of job completion time, throughput and job failure rate. Zhuozhao Li, Haiying Shen, Walter B. Ligon III, Jeffrey Denton |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2016 | On-Demand Bandwidth Pricing for Congestion Control in Core Switches in Cloud NetworksabstractThe cloud networks use switches to transfer inbound and outbound traffic through the data centers. Access of multiple tenants to the limited bandwidth capacity over the network switches increases the data traffic congestion in the network. The highly congested switches are vulnerable to get overloaded, and consequently slow down the flow of data traffic in the network. This paper proposes a nonlinear pricing policy for on-demand bandwidth allocation that jointly maximizes the total satisfaction of tenants and minimizes the congestion in the core switches. The optimal schedule is found through the best response strategy, in which each tenant updates its bandwidth allocation at each step based on the updated load-dependent predetermined nonlinear bandwidth pricing functions. The updated bandwidth allocations converge to the optimal bandwidth schedule that balances the load over the core switches. The performance of proposed pricing policy is evaluated under different scenarios. Abouzar Ghavami, Zhuozhao Li, Haiying Shen |
CLOUD | 2 |
| 2016 | Performance Measurement on Scale-Up and Scale-Out Hadoop with Remote and Local File SystemsabstractMapReduce is a popular computing model for parallel data processing on large-scale datasets, which can vary from gigabytes to terabytes and petabytes. Though Hadoop MapReduce normally uses Hadoop Distributed File System (HDFS) local file system, it can be configured to use a remote file system. Then, an interesting question is raised: for a given application, which is the best running platform among the different combinations of scale-up and scale-out Hadoop with remote and local file systems. However, there has been no previous research on how different types of applications (e.g., CPU-intensive, data-intensive) with different characteristics (e.g., input data size) can benefit from the different platforms. Thus, in this paper, we conduct a comprehensive performance measurement of different applications on scale-up and scaleout clusters configured with HDFS and a remote file system (i.e., OFS), respectively. We identify and study how different job characteristics (e.g., input data size, the number of file reads/writes, and the amount of computations) affect the performance of different applications on the different platforms. This study is expected to provide a guidance for users to choose the best platform to run different applications with different characteristics in the environment that provides both remote and local storage, such as HPC cluster. Zhuozhao Li, Haiying Shen |
CLOUD | 1 |
| 2016 | Comparing application performance on HPC-based Hadoop platforms with local storage and dedicated storageabstractMany high-performance computing (HPC) sites extend their clusters to support Hadoop MapReduce for a variety of applications. However, HPC cluster differs from Hadoop cluster on the configurations of storage resources. In the Hadoop Distributed File System (HDFS), data resides on the compute nodes, while in the HPC cluster, data is stored on separate nodes dedicated to storage. Dedicated storage offloads I/O load from the compute nodes and provides more powerful storage. Local storage provides better locality and avoids contention for shared storage resources. To gain an insight of the two platforms, in this paper, we investigate the performance and resource utilization of different types (i.e., I/O-intensive, data-intensive and CPU-intensive) of applications on the HPC-based Hadoop platforms with local storage and dedicated storage. We find that the I/O-intensive and data-intensive applications with large input file size can benefit more from the dedicated storage, while these applications with small input file size can benefit more from the local storage. CPU-intensive applications with a large number of small-size input files benefit more from the local storage, while these applications with large-size input files benefit approximately equally from the two platforms. We verify our findings by trace-driven experiments on different types of jobs from the Facebook synthesized trace. This work provides guidance on choosing the best platform to optimize the performance of different types of applications and reduce system overhead. Zhuozhao Li, Haiying Shen, Jeffrey Denton, Walter B. Ligon III |
IEEE BigData | 1 |
| 2016 | Game Theory-Based Nonlinear Bandwidth Pricing for Congestion Control in Cloud NetworksabstractIn the cloud, the network links are shared among tenants, which makes them easy to get fully congested (overloaded). Overloaded links degrade the performance of tenants' applications, and impose additional costs to the cloud provider. In this paper, we propose a nonlinear bandwidth pricing policy for congestion control in the cloud network. In order to maximize social welfare (i.e., maximize the total satisfaction of the tenants while minimizing the congestion over the link), the cloud provider uses the nonlinear pricing policy that increases the unit price with increment of bandwidth usage. Each tenant competes for bandwidth allocation to maximize its utility (i.e., both maximize its own individual satisfaction and minimize its bandwidth payment cost). We design a game between tenants and the cloud provider, and show that there exists a unique optimal bandwidth schedule (Nash equilibrium) that jointly maximizes the social welfare and the utility of each tenant at the same time. In order to find the optimal schedule, we use an asynchronous-based best response strategy, in which each tenant updates its optimal bandwidth allocation based on the updated bandwidth payment function from the cloud provider. We prove that the updated bandwidth allocations converge to the optimal bandwidth schedule. In our simulation study and real implementation, we verify the performance of our proposed pricing mechanism under different scenarios. Abouzar Ghavami, Zhuozhao Li, Haiying Shen |
CloudCom | 2 |
| 2016 | Goodbye to Fixed Bandwidth Reservation: Job Scheduling with Elastic Bandwidth Reservation in CloudsabstractThe shared nature of cloud network infrastructures causes unpredictable network performance, which may degrade the performance of these applications. Recently, several works propose to explicitly reserve the network bandwidth in the cloud with virtual network abstraction models, which pre-specify the network bandwidth between virtual machines (VMs) for a tenant job. However, the pre-specification fails to exploit the elastic feature of the bandwidth resource (i.e., more reserved bandwidth within no-elongation threshold bandwidth leads to shorter job execution time and vice versa) in job scheduling. It is difficult for ordinary tenants (without specialized network knowledge) to estimate the exact needed bandwidth. In this paper, we propose a new cloud job scheduler, in which each tenant only needs to specify job deadline and each job's reserved bandwidth is elastically determined by leveraging the elastic feature to maximize the total job rewards, which represent the worth of successful completion by deadlines. Finally, the scheduler tries to reduce the execution time of each job. It also jointly considers the computational capacity of VMs and reserved VM bandwidth in job scheduling. Using trace-driven and real cluster experiments, we show the efficiency and effectiveness of our job scheduler in comparison with other scheduling strategies. Haiying Shen, Lei Yu 0002, Liuhua Chen, Zhuozhao Li |
CloudCom | 4 |
| 2016 | Learning Network Graph of SIR Epidemic Cascades Using Minimal Hitting Set Based ApproachabstractWe consider learning the underlying graph structure of a network in which infection spreads based on the observations of node infection times. We give an algorithm based on minimal hitting set to learn the exact underlying graph structure and provide sufficient condition on number of cascades required (i.e. sample complexity) for reliable recovery, which is shown to be O(logn), where n is the number of nodes in the graph. We then analytically evaluate performance of minimal hitting set approach in learning the degree distribution and detecting leaf nodes of a graph and provide a sufficient condition for its sample complexity which is shown to be lower than that of learning the whole graph. We also generalize the exact graph estimation problem to the problem of estimating the graph within a certain distortion, measured by edit distance. We show that this edit distance based graph estimator has a lower sample complexity. Our experimental results based on both synthetic network topologies and a real-world network trace show that our algorithm achieves superior performance than a previously proposed algorithm based on maximum likelihood. Zhuozhao Li, Haiying Shen, Kang Chen 0002 |
ICCCN | 1 |
| 2016 | New Bandwidth Sharing and Pricing Policies to Achieve a Win-Win Situation for Cloud Provider and TenantsabstractFor predictable application performance or fairness in network sharing in clouds, many bandwidth allocation policies have been proposed. However, with these policies, tenants are not incentivized to use idle bandwidth or prevent link congestion, and may even take advantage of the policies to gain unfair bandwidth allocation. Increasing network utilization while avoiding congestion not only benefits cloud provider but also the tenants by improving application performance. In this paper, we propose a new pricing model that sets different unit prices for reserved bandwidth, the bandwidth on congested links and on uncongested links, and makes the unit price for congested links proportional to their congestion degrees. We use game theory model to analyze tenants' behaviors in our model and the current pricing models, which shows the effectiveness of our model in providing the incentives. With the pricing model, we propose a network sharing policy to achieve both min-guarantee and proportionality, while prevent tenants from earning unfair bandwidth. We further propose methods for each virtual machine to arrange its traffic to reduce its unsatisfied demand and maximize its utility, while increase network utilization. As a result, our solution creates a win-win situation, where tenants strive to increase their benefits in bandwidth sharing, which also concurrently increases the utilities of cloud provider and other tenants. Our simulation and trace-driven experimental results show the effectiveness of our solutions in creating the win-win situation. Haiying Shen, Zhuozhao Li |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2015 | Designing a Hybrid Scale-Up/Out Hadoop Architecture Based on Performance Measurements for High Application PerformanceabstractSince scale-up machines perform better for jobs with small and median (KB, MB) data sizes while scale-out machines perform better for jobs with large (GB, TB) data size, and a workload usually consists of jobs with different data size levels, we propose building a hybrid Hadoop architecture that includes both scale-up and scale-out machines, which however is not trivial. The first challenge is workload data storage. Thousands of small data size jobs in a workload may overload the limited local disks of scale-up machines. Jobs from scale-up and scale-out machines may both request the same set of data, which leads to data transmission between the machines. The second challenge is to automatically schedule jobs to either scale-up or scale-out cluster to achieve the best performance. We conduct a thorough performance measurement of different applications on scale-up and scale-out clusters, configured with Hadoop Distributed File System (HDFS) and a remote file system (i.e., OFS), respectively. We find that using OFS rather than HDFS can solve the data storage challenge. Also, we identify the factors that determine the performance differences on the scale-up and scale-out clusters and their cross points to make the choice. Accordingly, we design and implement the hybrid scale-up/out Hadoop architecture. Our trace-driven experimental results show that our hybrid architecture outperforms both the traditional Hadoop architecture with HDFS and with OFS in terms of job completion time. Zhuozhao Li, Haiying Shen |
ICPP | 1 |
| 2014 | New bandwidth sharing and pricing policies to achieve a win-win situation for cloud provider and tenantsabstractFor predictable application performance or fairness in network sharing in clouds, many bandwidth allocation policies have been proposed. However, with these policies, tenants are not incentivized to use idle bandwidth or prevent link congestion, and may even take advantage of the policies to gain unfair bandwidth allocation. Increasing network utilization while avoiding congestion not only benefits cloud provider but also the tenants by improving application performance. In this paper, we propose a new pricing model that sets different unit prices for reserved bandwidth, the bandwidth on congested links and on uncongested links, and makes the unit price for congested links proportional to their congestion degrees. We use game theory model to analyze tenants' behaviors in our model and the current pricing models, which shows the effectiveness of our model in providing the incentives. With the pricing model, we propose a network sharing policy to achieve both min-guarantee and proportionality, while prevent tenants from earning unfair bandwidth. We further propose methods for each virtual machine to arrange its traffic to maximize its utility. As a result, our solution creates a win-win situation, where tenants strive to increase their benefits in bandwidth sharing, which also concurrently increases the utilities of cloud provider and other tenants. Our simulation and trace-driven experimental results show the effectiveness of our solution in creating the win-win situation. Haiying Shen, Zhuozhao Li |
INFOCOM | 2 |