Ligang He

dblp:36/5655 · DBLP profile ↗
← Back
122ranked-venue papers
14as first author
51since 2021 · last 2026
0000-0002-5671-0576ORCID · conflict

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

Systems, architecture and hardware · 86 · 14 first-author · 32 since 2021Artificial intelligence and machine learning · 9 · 8 since 2021Software engineering, systems software and programming languages · 8 · 5 since 2021Computer networks · 7 · 4 since 2021Databases, data management, data science and information retrieval · 4 · 1 since 2021Graphics, computer vision, multimedia, augmented reality and games · 3 · 3 since 2021Applied, interdisciplinary, general and emerging computing · 3 · 1 since 2021Human-computer interaction and ubiquitous computing · 2Theory of computation · 2
YearPublicationVenuePosition
2026 HieraNTT: A Memory Hierarchy-Aware Data Access Architecture for Efficient Number Theoretic Transform on GPU
Qian Xiong, Weiliang Ma, Ligang He, Yufan Bai, Yao Chen 0008, Hai Jin 0001, Xuanhua Shi
APPT3
2026 TAGT: An Efficient Graph Transformer Accelerator with Topology-aware Sparsification and Merging
abstract
Graph Transformers (GTs) have emerged as a powerful paradigm for graph representation learning, as their attention mechanism can capture long-range dependencies and model complex structural interactions beyond the local messagepassing scope of conventional Graph Neural Networks (GNNs). This capability has enabled GTs to achieve strong accuracy across important domains, including recommendation systems and VLSI congestion prediction. However, the global attention mechanism in GTs requires each vertex to attend to all other vertices, incurring O(N2) computation and intermediate data movement. As graph size increases, this quadratic complexity leads to prohibitive computational overhead and excessive offchip memory traffic, fundamentally limiting the scalability and efficiency of GT execution. In this paper, we propose TAGT, the first efficient topologyaware Graph Transformer accelerator designed to mitigate these performance bottlenecks. Specifically, we integrate a topologyaware sparsification and merging approach into the accelerator design that dramatically reduces the O(N2) complexity. TAGT introduces a structure-aware sparse subgraph, termed the Topology Dependency Subgraph (TDS), which exploits inherent topological dependencies and reduces the number of attended edges to O(N log N) on average. The TDS is designed to retain local neighborhood structure while capturing essential higherorder interrelationships. By performing attention on the TDS, TAGT approximates global attention over the entire graph with negligible accuracy loss while eliminating most unnecessary computations and off-chip data movements. To fully harness the performance potential of this approach, TAGT incorporates a datadriven loading and merging engine to minimize off-chip memory accesses and reduce TDS construction overhead on the fly. TAGT also introduces a TDS-based fast attention unit to improve the parallelism of attention computation. We implement and evaluate TAGT on a Xilinx Alveo U280 FPGA card. Experimental results show that TAGT achieves average speedups of 175.4 × and 18.6 ×, together with energy savings of 217.2 × and 24.8 ×, over state-ofthe-art software GT solutions on Intel Xeon CPUs and NVIDIA A100 GPUs, respectively. Compared with representative GNN accelerators, including FlowGNN, MEGA, and BingoGCN, TAGT delivers average speedups of 8.2 ×, 6.9 ×, and 4.7 ×, and energy savings of 9.3 ×, 7.5 ×, and 5.2 ×, respectively.
Ligang He, Jin Zhao 0003
ISCA3
2026 DTMiner: A Data-Centric System for Efficient Temporal Motif Mining
abstract
Mining temporal motifs in temporal graphs is essential for many critical applications. Although several solutions have been proposed to handle temporal motif mining, they still suffer from substantial inefficiencies due to significant redundant graph traversals and fragmented memory access, both caused by irregular search tree expansions across different motif matching tasks. In this work, we observe that data accesses issued by these tasks exhibit strong spatial similarity and temporal monotonicity. Based on these observations, this paper proposes an efficient data-centric temporal motif mining system DTMiner, which introduces a novel Load-Explore-Synchronize (LES) execution model to efficiently regularize data accesses to the common temporal graph data among different tasks. Specifically, DTMiner enables the temporal graph chunks to be sequentially loaded into the cache in temporal order and then triggers all relevant tasks to explore only these loaded data for search tree expansions in a fine-grained synchronization mechanism. In this way, different tasks can share the graph traversal corresponding to the same chunks, while fragmented memory accesses are restricted to the graph data residing in the cache, significantly reducing data access overhead. Experimental results demonstrate that DTMiner achieves 1.14×-11.98× performance improvement in comparison with the state-of-the-art temporal motif mining solutions.
Yinbo Hou, Hao Qi 0004, Ligang He, Jin Zhao 0003, Yu Zhang 0027, Longlong Lin, Lin Gu 0002, Wenbin Jiang 0001, Xiaofei Liao, Hai Jin 0001
PPoPP3
2026 ACT-Agent: Affinity-Cross Transformer for Point Cloud Registration via Reinforcement Learning
abstract
ABSTRACT Point cloud registration, a core task in 3D computer vision for aligning two point clouds via rotation and translation, underpins critical applications like robotic navigation and 3D reconstruction. Classical methods (e.g., Iterative Closest Point) easily converge to local minima under poor initial alignment. Deep learning–based approaches, while efficient, suffer from high annotation costs for large‐scale data. Existing reinforcement learning (RL)‐based methods rely on simple PointNet feature extractors, which are insensitive to local geometric details and thus yield suboptimal registration precision. To address these challenges, we propose ACT‐Agent: Affinity‐Cross Transformer for point cloud registration via reinforcement learning, a novel method that formulates point cloud registration as an RL Markov decision process for iterative optimisation. We leverage Pointnet and Affinity‐Cross Transformer to extract and enhance expressive salient features and assign adaptive weights to channels based on their relative importance. We use RL to autonomously learn from feedback in the environment, freeing ourselves from dependence on data annotation. Experimental results on ModelNet40 (synthetic data) and ScanObjectNN (real‐world data) demonstrate that our proposed ACT‐Agent achieves higher accuracy, efficiency, and generalisation ability than the state‐of‐the‐art methods of point cloud registration.
Fengguang Xiong, Haixin Gong, Qiao Ma, Yingbo Jia, Ruize Guo, Ligang He, Liqun Kuang, Xie Han 0001
IET Image Process.7
2026 DeepPAT: Deep position-aware transformer with mixed dataset for robust point cloud registration
Fengguang Xiong, Ligang He, Liqun Kuang, Xie Han 0001
Neurocomputing3
2026 Interpreting networks via semantic probes with correction and counterfactuals
Keyang Cheng, Ligang He, Maozhen Li
Pattern Recognit.3
2025 TempGraph: An Efficient Chain-driven Temporal Graph Computing Framework on the GPU
abstract
Tackling temporal path problems in temporal graphs is essential for time-sensitive applications. Although many solutions have been proposed to handle temporal path problems, due to the intrinsic time constraints, these solutions require the vertices of the temporal graph to be sequentially handled along the time-dependent chains (i.e., the temporal dependencies between these vertices) to form the temporal path. This sequential temporal nature poses the challenges of poor parallelism and slow convergence speed, preventing existing solutions from fully leveraging the massive parallelism and high internal bandwidth of GPU to handle temporal path problems. To overcome these challenges, this paper proposes TempGraph, an efficient chain-driven GPU-based temporal graph computing framework. Specifically, it transforms the temporal graph into a set of disjoint time-dependent chains that can elegantly expose the temporal dependency between the vertices while facilitating the fast path exploration along these chains over GPU. Furthermore, TempGraph employs a novel Generate-Activate-Compute execution model to decouple the temporal dependency between different chains through maintaining a set of shortcuts for them, which enables multiple chains to be concurrently handled by massive GPU threads, achieving fast convergence speed and high parallelism on the GPU. Experiments on an A100 GPU show that TempGraph outperforms the state-of-the-art GPU-based solutions by 3.0-16.2×. Besides, TempGraph on an A100 GPU gains 33.9-368.9× speedups compared to the cutting-edge CPU-based system TeGraph on a 128-core CPU machine.
Jin Zhao 0003, Qian Wang 0002, Ligang He, Yu Zhang 0027, Sheng Di, Bingsheng He, Hao Qi 0004, Longlong Lin, Linchen Yu, Xiaofei Liao, Hai Jin 0001
ASPLOS (3)3
2025 OHMiner: An Overlap-centric System for Efficient Hypergraph Pattern Mining
abstract
Hypergraph Pattern Mining (HPM) aims to identify all the instances of user-interested subhypergraphs (patterns) in hypergraphs, which has been widely used in various applications. However, existing solutions either need significant enumeration overhead because they extend subhypergraphs at the granularity of vertices, or suffer from massive redundant computations because they often need to repeatedly fetch and process the same incident hyperedges for different vertices. This paper presents an overlap-centric system named OHMiner to efficiently support HPM. OHMiner proposes an overlap-centric execution model to determine the subhypergraphs isomorphism through computing and comparing overlaps among hyperedges using set operations. This model aims to efficiently handle the vertices that collectively share the same incident hyperedges. To automatically and precisely retrieve an arbitrary pattern's overlapping semantics without performing redundant set computations, OHMiner further proposes a redundancy-free compiler, which constructs an Overlap Intersection Graph (OIG) for the pattern, optimizes the OIG, and generates an overlap-centric execution plan to guide the procedure of HPM. Moreover, OHMiner designs an overlap-centric parallel execution engine, which adopts an incremental overlap-pruned approach to fast validate candidates for HPM. Additionally, it proposes a degree-aware data store to support efficient generation of candidates. Through evaluating OHMiner on a broad range of real-world hypergraphs with various patterns, our experimental results show that OHMiner outperforms the state-of-the-art HPM system by 5.4×-22.2×.
Hao Qi 0004, Ligang He, Yu Zhang 0027, Minzhi Cai, Jingxin Dai, Bingsheng He, Hai Jin 0001, Zhan Zhang 0003, Jin Zhao 0003, Hengshan Yue, Xiaofei Liao
EuroSys3
2025 TSGAN: A Lightweight Teacher-Student-Based GAN Framework for the Edge-Cloud Computing
abstract
Generative Adversarial Networks (GANs) have proven to be effective in generating synthesized data. However, their use often requires large amounts of data and sufficient computing power for model training. This presents significant challenges for GAN training in edge devices, which typically lack the necessary resources for GAN training. Additionally, the data collected by these devices often contain sensitive information that requires privacy preservation. The data volume may also be insufficient for GAN training, and the data collected by individual edge devices may have both common and unique features across the network. To address these challenges, we propose TSGAN, a lightweight GAN framework for Edge-Cloud computing. Coordinated by a Cloud server, TSGAN allows a network of resource-limited edge devices to train GAN models for privacy-preserving data generation. We propose a teacher-student model to enable edge devices to generate high-quality data. We also propose a novel deployment mechanism that facilitates effective distributed learning across edge devices and the Cloud server, while preventing the Cloud server generating the synthetic data from accessing the data collected by edge devices. Finally, we introduce a joint restraint learning function that enhances the effectiveness of learning unique features from data on individual edge devices. We conducted extensive experiments. The results have verified the effectiveness of TSGAN.
Weitong Liao, Hao Wu 0058, Zikun Zhang, Mohammed H. Alghamdi, Hammam M. AlGhamdi, Ligang He
HPCC6
2025 Octopus: Decentralized Workflow-granular Scheduling for Serverless Workflow
abstract
With the continuous development of Serverless Computing, Serverless applications composed of multiple finegrained functions have been widely applied in various fields of real life. As a pre-defined logical abstraction of Serverless applications, Serverless Workflow describes the dependencies and data flow between functions, and is the mainstream paradigm of modern Serverless Computing. However, our investigation shows that traditional Master-Worker-based, function-granular Serverless Workflow Management Systems are no longer suitable for the multi-function composition and unpredictable high concurrency characteristics of current Serverless Workflow. The seemingly insignificant scheduling overhead of Serverless Workflows has become a non-trivial factor affecting the execution efficiency and scalability of Serverless Workflows. Therefore, we proposed a workflow-granular management paradigm and decentralized control to address these challenges. Following these methodologies, we implement Octopus to enable efficient workflow scheduling and execution across different levels of concurrency and cluster scales. Experiments indicate that, in high-concurrency environments, Octopus achieves up to a 90× reduction in scheduling overhead and can enhance execution efficiency by 7.5×. As the cluster size increases, Octopus shows acceptable overhead and high scalability.
Keming Wang, Liaoliao Feng, Ligang He, Chenlin Huang, Tao Xie 0012
ICDCS3
2025 TaGNN: An Efficient Topology-aware Accelerator for High-performance Dynamic Graph Neural Network
abstract
Dynamic Graph Neural Networks (DGNNs) have become powerful tools for analyzing continuously evolving graph data, combining Graph Neural Network (GNN) models to extract structural information and Recurrent Neural Network (RNN) models to capture temporal semantics across snapshots. However, despite extensive research, existing DGNN solutions still face significant limitations, particularly low data parallelism caused by their snapshot-by-snapshot execution. This sequential paradigm exacerbates memory contention due to irregular, repeated vertex feature accesses and enforces strict temporal dependencies. In this paper, we propose TaGNN, an efficient topology-aware DGNN accelerator that addresses these performance bottlenecks. Specifically, we present a topology-aware concurrent execution approach into the accelerator design that calculates the final features of affected vertices while ensuring that unaffected vertices are loaded and computed only once per layer across multiple snapshots, maximizing data parallelism while minimizing memory usage. TaGNN employs a cache-friendly storage format that compactly organizes affected vertices across multiple snapshots by their timestamps and topological characteristics, reducing indexing overhead and enhancing data locality. In addition, TaGNN further proposes a similarity-aware cell skipping strategy to alleviate the stringent temporal data dependencies. It selectively reuses the RNN results from the previous snapshot to bypass RNN operations in the current snapshot when the output features of the GNN module across two consecutive snapshots are similar, achieving significant efficiency gains with minimal accuracy loss. We have implemented and assessed TaGNN on a Xilinx Alveo U280 FPGA card. Experimental results show that TaGNN achieves average speedups of 535.2x and 84.3x, and energy savings of 742.6x and 104.9x over state-of-the-art software DGNNs on Intel Xeon CPUs and NVIDIA A100 GPUs, respectively. Compared to leading DGNN accelerators (i.e., DGNN-Booster, E-DGCN, and Cambricon-DG), TaGNN delivers average speedups of 13.5x, 10.2x, and 6.5x, and energy savings of 15.9x, 11.7x, and 7.8x, respectively.
Yu Zhang 0027, Ligang He, Bing Peng, Jin Zhao 0003, Zixiao Wang 0005, Hao Qi 0004, Hai Jin 0001
SC3
2025 Multi-agent dual actor-critic framework for reinforcement learning navigation
Fengguang Xiong, Yaodan Zhang, Xinhe Kuang, Ligang He, Xie Han 0001
Appl. Intell.4
2025 Learning from imbalance: Cross-server power prediction in large data centers via domain adaptation regression
Ruichao Mo, Weiwei Lin 0001, Guozhi Liu, Haolin Liu 0001, Ligang He
Expert Syst. Appl.5
2025 Towards imbalanced regression over distributionally biased data: A fast static approach
Wentai Wu, Ligang He, Weiwei Lin 0001, Jinyi Long, Zhiquan Liu 0001, C. L. Philip Chen
Inf. Softw. Technol.2
2025 RuYi: Optimizing Burst Buffer Through Automated, Fine-Grained Process-to-BB Mapping
abstract
Current supercomputers use an SSD-based storage layer called Burst Buffer (BB) to provide I/O-intensive applications with accelerated storage access. However, efficiently utilizing this limited and expensive storage remains a critical issue, creating an urgent need for implementing Quality of Service (QoS) in BB. To address this, we propose RuYi, a QoS-aware method to provide applications with bandwidth guarantees in the BB file system. RuYi tackles two main issues. First, it quantitatively profiles available bandwidth resources in BB to ensure reliable QoS, a crucial aspect seldom studied in the literature. Second, RuYi offers fine-grained process-level QoS via an innovative process-to-BB mapping, maximizing resource utilization—something not achievable with conventional coarse-grained compute-to-BB mapping. We evaluated RuYi on a subsystem of the leading exascale supercomputer Sunway, consisting of 4,000 compute nodes and 200 BB nodes. The experimental results demonstrate that RuYi achieves an impressive end-to-end bandwidth control accuracy of 97%, while improving BB utilization by up to 116% compared to conventional coarse-grained compute-to-BB mapping.
Yusheng Hua, Xuanhua Shi, Ligang He, Teng Zhang 0001, Hai Jin 0001, Yong Chen 0001
IEEE Trans. Computers3
2025 PFed-NS: An Adaptive Personalized Federated Learning Scheme Through Neural Network Segmentation
abstract
Federated Learning (FL) is typically deployed in a client-server architecture, which makes the Edge-Cloud architecture an ideal backbone for FL. A significant challenge in this setup arises from the diverse data feature distributions across different edge locations (i.e., non-IID data). In response, Personalized Federated Learning (PFL) approaches have been developed. Network segmentation-based PFL is an important approach to achieving PFL, in which the training network is divided into a global segment for server aggregation and a local segment maintained client-side. Existing methods determine the segmentation before the training, and the segmentation remains fixed throughout the PFL training. However, our investigation reveals that model representations vary as PFL progresses and the fixed segmentation may not deliver best performance across various training settings. To address this, we propose PFed-NS, a PFL framework based on adaptive network segmentation. This adaptive segmentation technique is composed of two elements: a mechanism for assessing divergence of clients’ probability density functions constructed from network layers’ outputs, and a model for dynamically establishing divergence thresholds, beyond which server aggregation is deemed detrimental. Further optimization strategies are proposed to reduce the computation and communication costs incurred by divergence modeling. Moreover, we propose a divergence-based BN strategy to optimize BN performance for network segmentation-based PFL. Extensive experiments have been conducted to compare PFed-NS against recent PFL models. The results demonstrate its superiority in enhancing model accuracy and accelerating convergence.
Ligang He, Zhigao Zhang, Shenyuan Ren
IEEE Trans. Computers2
2024 CDA-GNN: A Chain-driven Accelerator for Efficient Asynchronous Graph Neural Network
abstract
Asynchronous Graph Neural Network (AGNN) has attracted much research attention because it enables faster convergence speed than the synchronous GNN. However, existing software/hardware solutions suffer from redundant computation overhead and excessive off-chip communications for AGNN due to irregular state propagations along the dependency chains between vertices. This paper proposes a chain-driven asynchronous accelerator, CDA-GNN, for efficient AGNN inference. Specifically, CDA-GNN proposes a chain-driven asynchronous execution approach into novel accelerator design to regularize the vertex state propagations for fewer redundant computations and off-chip communications and also designs a chain-aware data caching method to improve data locality for AGNN. We have implemented and evaluated CDA-GNN on a Xilinx Alveo U280 FPGA card. Compared with the cutting-edge software solutions (i.e., Dorylus and AMP) and hardware solutions (i.e., BlockGNN and FlowGNN), CDA-GNN improves the performance of AGNN inference by an average of 1,173x, 182.4x, 10.2x, and 7.9x and saves energy by 2,241x, 242.2x, 12.4x, and 8.9x, respectively.
Yu Zhang 0027, Ligang He, Donghao He, Qikun Li, Jin Zhao 0003, Xiaofei Liao, Hai Jin 0001, Lin Gu 0002, Haikun Liu
DAC3
2024 LSGraph: A Locality-centric High-performance Streaming Graph Engine
abstract
Streaming graph has been broadly employed across various application domains. It involves updating edges to the graph and then performing analytics on the updated graph. However, existing solutions either suffer from poor data locality and high computation complexity for streaming graph analytics, or need high overhead to search and move graph data to ensure ordered neighbors during streaming graph update.
Hao Qi 0004, Yiyang Wu, Ligang He, Yu Zhang 0027, Minzhi Cai, Hai Jin 0001, Zhan Zhang 0003, Jin Zhao 0003
EuroSys3
2024 RAHP: A Redundancy-aware Accelerator for High-performance Hypergraph Neural Network
abstract
Hypergraph Neural Network (HyperGNN) has emerged as a potent methodology for dissecting intricate multilateral connections among various entities. Current software/hardware solutions leverage a sequential execution model that relies on hyperedge and vertex indices for conducting standard matrix operations for HyperGNN inference. Yet, they are impeded by the dual challenges of redundant computation and irregular memory access overheads. This is primarily due to the frequent and repetitive access and updating of a number of feature vectors corresponding to the same hyperedges and vertices. To address these challenges, we propose the first redundancy-aware accelerator, RAHP, which enables high performance execution of HyperGNN inference. Specifically, we present a redundancy-aware asynchronous execution approach into the accelerator design for HyperGNN to reduce redundant computations and off-chip memory accesses. To unveil opportunities for data reuse and unlock the parallelism that existing HyperGNN solutions fail to capture, it prioritizes vertices with the highest degree as roots, prefetching other vertices along the hypergraph structure to capture the common vertices among multiple hyperedges, and synchronizing the computations of hyperedges and vertices in real-time. By such means, this facilitates the concurrent processing of relevant hyperedge and vertex computations of the common vertices along the hypergraph topology, resulting in smaller redundant computations overhead. Furthermore, by efficiently caching intermediate results of the common vertices, it curtails memory traffic and off-chip communications. To fully harness the performance potential of our proposed approach in the accelerator, RAHP incorporates a topology-driven data loading mechanism to minimize off-chip memory accesses on the fly. It is also endowed with an adaptive data synchronization scheme to mitigate the effects of conflicting updates of both hyperedges and vertices. Moreover, RAHP employs the similarity-based data caching strategy to further mitigate the overhead of redundant data transfers. We have implemented and assessed RAHP on a Xilinx Alveo U280 FPGA card. Experimental evaluations demonstrate that RAHP achieves average speedups of 439.2x and 64.7x for HyperGNN inference, alongside average energy savings of 542.8x and 84.2x, compared to the cutting-edge software-based HyperGNN implementations on Intel Xeon CPUs and NVIDIA A100 GPUs, respectively. Additionally, in the realm of HyperGNN inference, RAHP secures average speedups of 7.8x, 5.4x, and 3.8x, and average energy savings of 10.2x, 8.9x, and 6.5x over the foremost GNN accelerators, i.e., FlowGNN, LL-GNN, and ReGNN, respectively.
Yu Zhang 0027, Ligang He, Yingqi Zhao, Xintao Li, Ruida Xin, Jin Zhao 0003, Xiaofei Liao, Haikun Liu, Bingsheng He, Hai Jin 0001
MICRO3
2024 SARAD: Spatial Association-Aware Anomaly Detection and Diagnosis for Multivariate Time Series
abstract
Anomaly detection in time series data is fundamental to the design, deployment, and evaluation of industrial control systems. Temporal modeling has been the natural focus of anomaly detection approaches for time series data. However, the focus on temporal modeling can obscure or dilute the spatial information that can be used to capture complex interactions in multivariate time series. In this paper, we propose SARAD, an approach that leverages spatial information beyond data autoencoding errors to improve the detection and diagnosis of anomalies. SARAD trains a Transformer to learn the spatial associations, the pairwise inter-feature relationships which ubiquitously characterize such feedback-controlled systems. As new associations form and old ones dissolve, SARAD applies subseries division to capture their changes over time. Anomalies exhibit association descending patterns, a key phenomenon we exclusively observe and attribute to the disruptive nature of anomalies detaching anomalous features from others. To exploit the phenomenon and yet dismiss non-anomalous descent, SARAD performs anomaly detection via autoencoding in the association space. We present experimental results to demonstrate that SARAD achieves state-of-the-art performance, providing robust anomaly detection and a nuanced understanding of anomalous events.
Zhihao Dai, Ligang He, Shuang-Hua Yang, Matthew Leeke
NeurIPS2
2024 A Data-Distillation-Enhanced Autoencoder for Detecting Anomalous Gas Consumption
abstract
The number of natural gas users has been growing rapidly in China due to the promotion of clean energy and the economic benefits of natural gas, especially in businesses and industries. Though the infrastructures for gas supplies have been highly improved, gas providers are still suffering from various problems such as malfunctioning gas meters, gas leakage, gas theft etc. With the development of the Internet of Things, smart gas meters have been widely adopted by gas providers to collect real-time gas consumption data for billing purposes which can also serve as a basis for anomaly detection. One challenge of using such data for anomaly detection is that it is difficult to obtain sufficient labelled data for model training. To address this challenge, we propose DAE, a data distillation enhanced autoencoder for detecting anomalous gas consumption, which consists of three modules. The first module preprocesses the raw meter readings and carries out a rule-based anomaly detection. The second module extracts the normal gas usage patterns via an integration of correlation and clustering based consistency evaluation methods. The extracted normal usage patterns are then used in the third module to train an autoencoder for anomaly detection. DAE intends to provide a method to detect anomalous gas consumption induced by various causes such that manual inspection can be largely reduced. Moreover, DAE does not require user-specific information and can be applied to different types of gas users. Based on a real-world gas consumption dataset, we carry out a set of experiments and show that DAE outperforms the existing and improves the F1 score by an average of 7.4% for restaurant users and 5.7% for canteen users.
Yujue Zhou, Jie Jiang 0011, Shuang-Hua Yang, Ligang He, Guozhong Zhu, Yali Qing
IEEE Internet Things J.4
2024 MMDataLoader: Reusing Preprocessed Data Among Concurrent Model Training Tasks
abstract
Data preprocessing plays an important role in deep learning, which directly affects the training efficiency. Data preprocessing is performed on the CPU. The preprocessed data are then fed to the models that are trained on the GPU. We observe that data preprocessing on the CPU can potentially create a bottleneck in the entire process of a model training task. In order to tackle this issue, we have developed MMDataLoader, which enables reusing preprocessed data among multiple model training tasks. MMDataLoader automatically constructs a data preprocessing pipeline based on each task's specific preprocessing workflow, allowing for maximum data reuse and reduced computing workload on the CPU. Unlike conventional data loaders that operate at the task level and provide data provision services to specific training tasks, MMDataLoader operates at the server level and provides data for all concurrently running tasks. We have conducted extensive experiments. The results show that MMDataLoader can significantly increase preprocessing throughput without affecting model convergence when compared to conventional methods where model training tasks are executed concurrently. For instance, with three tasks running, the preprocessing throughput can increase by 1.6x to 3.15x, depending on the tasks being executed and the proportion of preprocessing operations that are shared among them.
Hai Jin 0001, Zhanyang Zhu, Ligang He, Yusheng Hua, Xuanhua Shi
IEEE Trans. Computers3
2023 PSMiner: A Pattern-Aware Accelerator for High-Performance Streaming Graph Pattern Mining
abstract
Streaming Graph Pattern Mining (GPM) has been widely used in many application fields. However, the existing streaming GPM solution suffers from many unnecessary explorations and isomorphism tests, while the existing static GPM ones require many repetitive operations to compute the full graph. In this paper, we propose a pattern-aware incremental execution approach and design the first streaming GPM accelerator called PSMiner, which integrates multiple optimizations to reduce redundant computation and improve computing efficiency. We have conducted extensive experiments. The results show that compared with the state-of-the-art software and hardware solutions, PSMiner achieves the average speedups of 770.9× and 60.4×, respectively.
Hao Qi 0004, Yu Zhang 0027, Ligang He, Haoyu Lu, Jin Zhao 0003, Hai Jin 0001
DAC3
2023 Service-Aware Cooperative Task Offloading and Scheduling in Multi-access Edge Computing Empowered IoT
Ming Tao 0001, Xueqiang Li 0001, Ligang He
ICA3PP (2)4
2023 Multi-view stereo network with point attention
Zhuoer Gu, Xie Han 0001, Ligang He, Fusheng Sun, Shichao Jiao
Appl. Intell.4
2023 Self-Supervised Leaf Segmentation under Complex Lighting Conditions
Xufeng Lin, Chang-Tsun Li, Scott D. Adams, Abbas Z. Kouzani, Richard Jiang 0001, Ligang He, Yongjian Hu, Michael Vernon, Egan H. Doeven, Lawrence Webb, Todd Mcclellan, Adam Guskic
Pattern Recognit.6
2023 GraphTune: An Efficient Dependency-Aware Substrate to Alleviate Irregularity in Concurrent Graph Processing
abstract
With the increasing need for graph analysis, massive Concurrent iterative Graph Processing (CGP) jobs are usually performed on the common large-scale real-world graph. Although several solutions have been proposed, these CGP jobs are not coordinated with the consideration of the inherent dependencies in graph data driven by graph topology. As a result, they suffer from redundant and fragmented accesses of the same underlying graph dispersed over distributed platform, because the same graph is typically irregularly traversed by these jobs along different paths at the same time. In this work, we develop GraphTune , which can be integrated into existing distributed graph processing systems, such as D-Galois, Gemini, PowerGraph, and Chaos, to efficiently perform CGP jobs and enhance system throughput. The key component of GraphTune is a dependency-aware synchronous execution engine in conjunction with several optimization strategies based on the constructed cross-iteration dependency graph of chunks. Specifically, GraphTune transparently regularizes the processing behavior of the CGP jobs in a novel synchronous way and assigns the chunks of graph data to be handled by them based on the topological order of the dependency graph so as to maximize the performance. In this way, it can transform the irregular accesses of the chunks into more regular ones so that as many CGP jobs as possible can fully share the data accesses to the common graph. Meanwhile, it also efficiently synchronizes the communications launched by different CGP jobs based on the dependency graph to minimize the communication cost. We integrate it into four cutting-edge distributed graph processing systems and a popular out-of-core graph processing system to demonstrate the efficiency of GraphTune. Experimental results show that GraphTune improves the throughput of CGP jobs by 3.1∼6.2, 3.8∼8.5, 3.5∼10.8, 4.3∼12.4, and 3.8∼6.9 times over D-Galois, Gemini, PowerGraph, Chaos, and GraphChi, respectively.
Jin Zhao 0003, Yu Zhang 0027, Ligang He, Qikun Li, Xiaofei Liao, Hai Jin 0001, Lin Gu 0002, Haikun Liu, Bingsheng He, Ji Zhang 0001, Xianzheng Song, Lin Wang 0098, Jun Zhou 0011
ACM Trans. Archit. Code Optim.3
2023 Waterwave: A GPU Memory Flow Engine for Concurrent DNN Training
abstract
Training Deep Neural Networks (DNN) concurrently is becoming increasingly important for deep learning practitioners, e.g.,hyperparameter optimization (HPO)andneural architecture search (NAS). The GPU memory capacity is the impediment that prohibits multiple DNNs from being trained on the same GPU due to the large memory usage during training. In this paper, we proposeWaterwave, a GPU memory flow engine for concurrent deep learning training. First, to address the memory explosion brought by the long time lag between memory allocation and deallocation time, we develop an allocator tailored for multi-streams. By making the allocator aware of the stream information, aprioritized allocationis conducted based on the chunk'ssynchronizationattributes, allowing us to provide useable memory after scheduling rather than waiting it to be really released after GPU computation. Second,Waterwavepartitions the compute graph to a set of continuousnode groupsand then performs finer-grained scheduling:NodeGroup pipeline execution, to guarantee a proper memory requests order.Waterwavecan accomplish up to 96.8% of the maximum batch size of solo training. Additionally, in scenarios with high memory demand,Waterwavecan outperform existing spatial sharing and temporal sharing by up to 12x and 1.49x, respectively.
Xuanhua Shi, Ligang He, Yunfei Zhao 0001, Hai Jin 0001
IEEE Trans. Computers3
2023 TurboGNN: Improving the End-to-End Performance for Sampling-Based GNN Training on GPUs
abstract
Graph Neural Networks(GNN) have evolved as powerful models for graph representation learning. Sampling-based training methods have been introduced to train large graphs without compromising accuracy. However, it is challenging for the existing GNN systems to effectively utilize multi-core accelerators, especially GPUs, due to a large number of atomic operations and unbalanced workload originating from the serial execution of multiple GNN processing stages. In this paper, we propose a combination of optimization techniques to accelerate the end-to-end performance of the sampling-based GNN training process. Specifically, we propose an adaptive share memory-based sampling technique and a degree-guided thread block scheduling strategy to optimize the graph sampling. Further, based on the observations of resource demand in different training stages, we propose an asynchronous pipeline-based scheduling method, which accelerates the GNN training by decoupling different training stages into a pipeline and therefore improves the GPU resource utilization significantly. The experimental results show that compared with the existing work, the proposed methods can achieve up to 5.6X performance speedup in the end-to-end performance.
Xuanhua Shi, Ligang He, Hai Jin 0001
IEEE Trans. Computers3
2023 FedProf: Selective Federated Learning Based on Distributional Representation Profiling
abstract
Federated Learning (FL) has shown great potential as a privacy-preserving solution to learning from decentralized data that are only accessible to end devices (i.e., clients). The data locality constraint offers strong privacy protection but also makes FL sensitive to the condition of local data. Apart from statistical heterogeneity, a large proportion of the clients, in many scenarios, are probably in possession of low-quality data that are biased, noisy or even irrelevant. As a result, they could significantly slow down the convergence of the global model we aim to build and also compromise its quality. In light of this, we first present a new view of local data by looking into the representation space and observing that they converge in distribution to Normal distributions before activation. We provide theoretical analysis to support our finding. Further, we proposeFedProf, a novel algorithm for optimizing FL over non-IID data of mixed quality. The key of our approach is a distributional representation profiling and matching scheme that uses the global model to dynamically profile data representations and allows for low-cost, lightweight representation matching. Using the scheme we sample clients adaptively in FL to mitigate the impact of low-quality data on the training process. We evaluated our solution with extensive experiments on different tasks and data conditions under various FL settings. The results demonstrate that the selective behavior of our algorithm leads to a significant reduction in the number of communication rounds and the amount of time (up to 2.4× speedup) for the global model to converge and also provides accuracy gain.
Wentai Wu, Ligang He, Weiwei Lin 0001, Carsten Maple
IEEE Trans. Parallel Distributed Syst.2
2023 TurboMGNN: Improving Concurrent GNN Training Tasks on GPU With Fine-Grained Kernel Fusion
abstract
Graph Neural Networks(GNN) have evolved as powerful models for graph representation learning. Many works have been proposed to support GNN training efficiently on GPU. However, these works only focus on a single GNN training task such as operator optimization, task scheduling, and programming model. Concurrent GNN training, which is needed in the applications such as neural network structure search, has not been explored yet. This work aims to improve the training efficiency of the concurrent GNN training tasks on GPU by developing fine-grained methods to fuse the kernels from different tasks. Specifically, we propose a fine-grainedSparse Matrix Multiplication(SpMM) based kernel fusion method to eliminate redundant accesses to graph data. In order to increase the fusion opportunity and reduce the synchronization cost, we further propose a novel technique to enable the fusion of the kernels in forward and backward propagation. Finally, in order to reduce the resource contention caused by the increased number of concurrent, heterogeneous GNN training tasks, we propose an adaptive strategy to group the tasks and match their operators according to resource contention. We have conducted extensive experiments, including kernel- and model-level benchmarks. The results show that the proposed methods can achieve up to 2.6X performance speedup.
Xuanhua Shi, Ligang He, Hai Jin 0001
IEEE Trans. Parallel Distributed Syst.3
2022 TDGraph: a topology-driven accelerator for high-performance streaming graph processing
abstract
Many solutions have been recently proposed to support the processing of streaming graphs. However, for the processing of each graph snapshot of a streaming graph, the new states of the vertices affected by the graph updates are propagated irregularly along the graph topology. Despite the years' research efforts, existing approaches still suffer from the serious problems of redundant computation overhead and irregular memory access, which severely underutilizes a many-core processor. To address these issues, this paper proposes a topology-driven programmable accelerator TDGraph, which is the first accelerator to augment the many-core processors to achieve high performance processing of streaming graphs. Specifically, we propose an efficient topology-driven incremental execution approach into the accelerator design for more regular state propagation and better data locality. TDGraph takes the vertices affected by graph updates as the roots to prefetch other vertices along the graph topology and synchronizes the incremental computations of them on the fly. In this way, most state propagations originated from multiple vertices affected by different graph updates can be conducted together along the graph topology, which help reduce the redundant computations and data access cost. Besides, through the efficient coalescing of the accesses to vertex states, TDGraph further improves the utilization of the cache and memory bandwidth. We have evaluated TDGraph on a simulated 64-core processor. The results show that, the state-of-the-art software system achieves the speedup of 7.1~21.4 times after integrating with TDGraph, while incurring only 0.73% area cost. Compared with four cutting-edge accelerators, i.e., HATS, Minnow, PHI, and DepGraph, TDGraph gains the speedups of 4.6~12.7, 3.2~8.6, 3.8~9.7, and 2.3~6.1 times, respectively.
Jin Zhao 0003, Yun Yang 0001, Yu Zhang 0027, Xiaofei Liao, Lin Gu 0002, Ligang He, Bingsheng He, Hai Jin 0001, Haikun Liu
ISCA6
2022 Reveal training performance mystery between TensorFlow and PyTorch in the single GPU environment
Hulin Dai, Xuanhua Shi, Ligang He, Qian Xiong, Hai Jin 0001
Sci. China Inf. Sci.4
2022 Reinforcement Learning for Security-Aware Computation Offloading in Satellite Networks
abstract
The rise ofNewSpaceprovides a platform for small and medium businesses to commercially launch and operate satellites in space. In contrast to traditional satellites,NewSpaceprovides the opportunity for delivering computing platforms in space. However, computational resources within space are usually expensive and satellites may not be able to compute all computational tasks locally. Computation offloading (CO), a popular practice in Edge/Fog computing, could prove effective in saving energy and time in this resource-limited space ecosystem. However, CO alters the threat and risk profile of the system. In this article, we analyze security issues in space systems and propose a security-aware algorithm for CO. Our method is based on the reinforcement learning technique, deep deterministic policy gradient (DDPG). We show, using Monte-Carlo simulations, that our algorithm is effective under a variety of environment and network conditions and provide novel insights into the challenge of optimized location of computation.
Saurav Sthapit, Subhash Lakshminarayana, Ligang He, Gregory Epiphaniou, Carsten Maple
IEEE Internet Things J.3
2022 Deep cross-modal discriminant adversarial learning for zero-shot sketch-based image retrieval
Shichao Jiao, Xie Han 0001, Fengguang Xiong, Xiaowen Yang, Huiyan Han, Ligang He, Liqun Kuang
Neural Comput. Appl.6
2022 A Structure-Aware Storage Optimization for Out-of-Core Concurrent Graph Processing
abstract
With the huge demand for graph analytics in many real-world applications, massive iterative graph processing jobs are concurrently performed on the same graphs and suffer from significant high data access cost. To lower the data access cost toward high performance, several out-of-core concurrent graph processing solutions are recently designed to handle concurrent jobs by enabling these jobs to share the accesses of the same graph data. However, the set of active vertices in each partition are usually different for various concurrent jobs and also evolve with time, where some high-degree ones (or calledhub-vertices) of these active vertices require more iterations to converge due to the power-law property of real-world graphs. In consequence, existing solutions still suffer from much unnecessary I/O traffic, because they have to entirely load each partition into the memory for concurrent jobs even if most vertices in this partition are inactive and may be shared by a few jobs. In this paper, we propose an efficient structure-aware storage system, called GraphSO, for higher throughput of the execution of concurrent graph processing jobs. It can be integrated into existing out-of-core graph processing systems to promote the execution efficiency of concurrent jobs with lower I/O overhead. The key design of GraphSO is a fine-grained storage management scheme. Specifically, it logically divides the partitions of existing graph processing systems into a series of small same-sized chunks. At runtime, these small chunks with active vertices are judiciously loaded by GraphSO to construct new logical partitions (i.e., each logical partition is a subset of active chunks) for existing graph processing systems to handle, where the most-frequently-used chunks are preferentially loaded to construct the logical partitions and the other ones are delayed to wait to be required by more jobs. In this way, it can effectively spare the cost of loading the graph data associated with the inactive vertices with low repartitioning overhead and can also enable the loaded graph data to be fully shared by concurrent jobs. Moreover, GraphSO also designs a buffering strategy to efficiently cache the most-frequently-used chunks in the main memory to further minimize the I/O traffic by avoiding repeated load of them. Experimental results show that GraphSO improves the throughput of GridGraph, GraphChi, X-Stream, DynamicShards, LUMOS, Graphene, and Wonderland by 1.4-3.5 times, 2.1-4.3 times, 1.9-4.1 times, 1.9-2.9 times, 1.5-3.1 times, 1.3-1.5 times, and 1.3-2.7 times after integrating with them, respectively.
Xiaofei Liao, Jin Zhao 0003, Yu Zhang 0027, Bingsheng He, Ligang He, Hai Jin 0001, Lin Gu 0002
IEEE Trans. Computers5
2022 PSNet: Fast Data Structuring for Hierarchical Deep Learning on Point Cloud
abstract
In order to retain more feature information of local areas on a point cloud, local grouping and subsampling are the necessary data structuring steps in most hierarchical deep learning models. Due to the disorder nature of the points in a point cloud, the significant time cost may be consumed when grouping and subsampling the points, which consequently results in poor scalability. This paper proposes a fast data structuring method called PSNet (Point Structuring Net). PSNet transforms the spatial features of the points and matches them to the features of local areas in a point cloud. PSNet achieves grouping and sampling at the same time while the existing methods process sampling and grouping in two separate steps (such as using FPS plus kNN). PSNet performs feature transformation pointwise while the existing methods uses the spatial relationship among the points as the reference for grouping. Thanks to these features, PSNet has two important advantages: 1) the grouping and sampling results obtained by PSNet is stable and permutation invariant; and 2) PSNet can be easily parallelized. PSNet can replace the data structuring methods in the mainstream point cloud deep learning models in a plug-and-play manner. We have conducted extensive experiments. The results show that PSNet can improve the training and inference speed significantly while maintaining the model accuracy.
Luyang Li 0001, Ligang He, Jinjin Gao, Xie Han 0001
IEEE Trans. Circuits Syst. Video Technol.2
2022 Developing an Unsupervised Real-Time Anomaly Detection Scheme for Time Series With Multi-Seasonality
abstract
On-line detection of anomalies in time series is a key technique used in various event-sensitive scenarios such as robotic system monitoring, smart sensor networks and data center security. However, the increasing diversity of data sources and the variety of demands make this task more challenging than ever. First, the rapid increase in unlabeled data means supervised learning is becoming less suitable in many cases. Second, a large portion of time series data have complex seasonality features. Third, on-line anomaly detection needs to be fast and reliable. In light of this, we have developed a prediction-driven, unsupervised anomaly detection scheme, which adopts a backbone model combining the decomposition and the inference of time series data. Further, we propose a novel metric, Local Trend Inconsistency (LTI), and an efficient detection algorithm that computes LTI in a real-time manner and scores each data point robustly in terms of its probability of being anomalous. We have conducted extensive experimentation to evaluate our algorithm with several datasets from both public repositories and production environments. The experimental results show that our scheme outperforms existing representative anomaly detection algorithms in terms of the commonly used metric, Area Under Curve (AUC), while achieving the desired efficiency.
Wentai Wu, Ligang He, Weiwei Lin 0001, Yuhua Cui, Carsten Maple, Stephen A. Jarvis
IEEE Trans. Knowl. Data Eng.2
2022 LoomIO: Object-Level Coordination in Distributed File Systems
abstract
Device-level interference is recognized as a major cause of the performance degradation in distributed file systems. Although the approaches of mitigating interference through coordination at application-level, middleware-level, and server-level have shown beneficial results in previous studies, we find their effectiveness is largely reduced since I/O requests are re-arranged by underlying object file systems. In this research study, we prove that object-level coordination is critical and often the key to address the interference issue, as the scheduling of object requests determines the device-level accesses and thus determines the actual I/O bandwidth and latency. This article proposes an object-level coordination system, LoomIO, which uses an OBOP (One-Broadcast-One-Propagate) method and a time-limited coordination process to deliver highly efficient coordination service. Specifically, LoomIO enables object requests to achieve an optimized scheduling decision within a few milliseconds and largely mitigates the device-level interference. We have implemented a LoomIO prototye and integrated it into Ceph file system. The evaluation results show that LoomIO achieved the considerable improvements in resource utilization (by up to 35%), in I/O throughput (by up to 31%), and in 99th percentile latency (by up to 54%) compared to the K-optimal method which uses the same scheduling algorithm as LoomIO but does not have the coordination support.
Yusheng Hua, Xuanhua Shi, Hai Jin 0001, Wei Xie 0017, Ligang He, Yong Chen 0001
IEEE Trans. Parallel Distributed Syst.6
2022 An On-Line Virtual Machine Consolidation Strategy for Dual Improvement in Performance and Energy Conservation of Server Clusters in Cloud Data Centers
abstract
As data centers are consuming massive amount of energy, improving the energy efficiency of cloud computing has emerged as a focus of research. However, it is challenging to reduce energy consumption while maintaining system performance without increasing the risk of Service Level Agreement violations. Most of the existing consolidation approaches for virtual machines (VMs) consider system performance and Quality of Service (QoS) metrics as constraints, which usually results in large scheduling overhead and impossibility to achieve effective improvement in energy efficiency without sacrificing some system performance and cloud service quality. In this article, we first define the metrics of peak power efficiency and optimal utilization for heterogeneous physical machines (PMs). Then we propose Peak Efficiency Aware Scheduling (PEAS), a novel strategy of VM placement and reallocation for achieving dual improvement in performance and energy conservation from the perspective of server clusters. PEAS allocates and reallocates VMs in an on-line manner and always attempts to maintain PMs working in their peak power efficiency via VM consolidation. Extensive experiments on Cloudsim show that PEAS outperforms several energy-aware consolidation algorithms with regard to energy consumption, system performance as well as multiple QoS metrics.
Weiwei Lin 0001, Wentai Wu, Ligang He
IEEE Trans. Serv. Comput.3
2022 vChecker: an application-level demand-based co-scheduler for improving the performance of parallel jobs in Xen
abstract
Abstract Big data analysis requires the speedup of parallel computing. However, in the virtualized systems, the power of parallel computing is not fully exploited due to the limit of current VMM schedulers. Xen, one of the most popular virtualization platforms, has been widely used by industry to host parallel job. In practice, the virtualized systems are expected to accommodate both parallel jobs and serial jobs, and resource contention between virtual machines results in severe performance degradation of the parallel jobs. Moreover, the physical resource is vastly wasted during the communication process due to the ineffective scheduling of parallel jobs. Unfortunately, the existing schedulers of Xen are initially targeting at serial jobs, which are not capable of correctly scheduling the parallel jobs. This paper presents vChecker, an application-level co-scheduler which mitigates the performance degradation of the parallel job and optimizes the utilization of the hardware resource. Our co-scheduler takes number of available CPU cores in one hand, and satisfies need of the parallel jobs in other hand, which helps the credit scheduler of Xen to appropriately schedule the parallel job. As our co-scheduler is implemented at application level, no modifications on the hypervisor is required. The experimental result shows that the vChecker optimizes the performance of the parallel job in Xen and enhances the utilization of the system.
Ligang He, Shenyuan Ren, Rui Mao 0001
Wirel. Networks2
2022 Contention-aware prediction for performance impact of task co-running in multicore computers
abstract
Abstract In this paper, we investigate the influential factors that impact on the performance when the tasks are co-running on a multicore computers. Further, we propose the machine learning-based prediction framework to predict the performance of the co-running tasks. In particular, two prediction frameworks are developed for two types of task in our model: repetitive tasks (i.e., the tasks that arrive at the system repetitively) and new tasks (i.e., the task that are submitted to the system the first time). The difference between which is that we have the historical running information of the repetitive tasks while we do not have the prior knowledge about new tasks. Given the limited information of the new tasks, an online prediction framework is developed to predict the performance of co-running new tasks by sampling the performance events on the fly for a short period and then feeding the sampled results to the prediction framework. We conducted extensive experiments with the SPEC2006 benchmark suite to compare the effectiveness of different machine learning methods considered in this paper. The results show that our prediction model can achieve the accuracy of 99.38% and 87.18% for repetitive tasks and new tasks, respectively.
Shenyuan Ren, Ligang He, Chang-Tsun Li
Wirel. Networks2
2021 DepGraph: A Dependency-Driven Accelerator for Efficient Iterative Graph Processing
abstract
Many graph processing systems have been recently developed for many-core processors. However, for iterative graph processing, due to the dependencies between vertices' states, the propagations of new states of vertices are inherently conducted along graph paths sequentially and are also dependent on each other. Despite the years' research effort, existing solutions still severely underutilize many-core processors to quickly propagate the new states of vertices, suffering from slow convergence speed. In this paper, we propose a dependency-driven programmable accelerator, DepGraph, which couples with the core architecture of the many-core processor and can fundamentally alleviate the challenge of dependencies for faster state propagation. Specifically, we propose an effective dependency-driven asynchronous execution approach into novel microarchitecture designs for faster state propagations. DepGraph prefetches the vertices for the core on-the-fly along the dependency chains between their states and the active vertices' new states, aiming to effectively accelerate the propagations of the active vertices' new states and also ensure better data locality. Through transforming the dependency chains along the frequently-used paths into direct ones at runtime and maintaining these calculated direct dependencies as a set of fast shortcuts, called hub index, DepGraph further accelerates most state propagations. Also, many propagations do not need to wait for the completion of other propagations, which enables more propagations to be effectively conducted along the paths with higher degree of parallelism. The experimental results show that for iterative graph processing on a simulated 64-core processor, a cutting-edge software graph processing system can achieve 5.0-22.7 times speedup after integrating with our DepGraph while incurring only 0.6% area cost. In comparison with three state-of-the-art hardware solutions, i.e., HATS, Minnow, and PHI, DepGraph improves the performance by up to 3.0-14.2, 2.2-5.8, and 2.4-10.1 times, respectively.
Yu Zhang 0027, Xiaofei Liao, Hai Jin 0001, Ligang He, Bingsheng He, Haikun Liu, Lin Gu 0002
HPCA4
2021 Multi-Camera Logical Topology Inference via Conditional Probability Graph Convolution Network
abstract
In order to improve the efficiency of pedestrian retrieval and re-identification with numerous surveillance cameras, a novel multi-camera dynamic logical topology inference method is proposed, which includes:(1) A conditional probability graph convolution network(CPG) is designed, which samples and aggregates information of multi-order neighbor nodes according to the conditional probability. The CPG is employed to aggregate the influence of all nodes relative to the target node and calculate the global correlation between each node.(2) A dynamic spatio-temporal information aggregation model(STIA) in a multi-camera system is proposed. The dynamic logical topology of the multi-camera system is inferred based on the pedestrian’s walking direction and the spatio-temporal factors.(3) A novel correlation indicator in multi-camera system is proposed. This indicator detects and quantifies temporal and causal relationships within and across camera views by a designed time-delayed Jensen-Shannon divergence(TDJS). It can be used to measure the causal correlation between two camera nodes over a long time delay. Some ablation studies and simulation experiments are performed on a dataset collected from real scenes show that our method can efficiently infer the logical topology of multiple cameras.
Keyang Cheng, Qing Liu 0015, Rabia Tahir, Lubamba Kasangu Eric, Ligang He
ICME5
2021 LCCG: a locality-centric hardware accelerator for high throughput of concurrent graph processing
abstract
In modern data centers, massive concurrent graph processing jobs are being processed on large graphs. However, existing hardware/-software solutions suffer from irregular graph traversal and intense resource contention. In this paper, we propose LCCG, a Locality-Centric programmable accelerator that augments the many-core processor for achieving higher throughput of Concurrent Graph processing jobs. Specifically, we develop a novel topology-aware execution approach into the accelerator design to regularize the graph traversals for multiple jobs on-the-fly according to the graph topology, which is able to fully consolidate the graph data accesses from concurrent jobs. By reusing the same graph data among more jobs and coalescing the accesses of the vertices' states for these jobs, LCCG can improve the core utilization. We conduct extensive experiments on a simulated 64-core processor. The results show that LCCG improves the throughput of the cutting-edge software system by 11.3~23.9 times with only 0.5% additional area cost. Moreover, LCCG gains the speedups of 4.7~10.3, 5.5~13.2, and 3.8~8.4 times over state-of-the-art hardware graph processing accelerators (namely, HATS, Minnow, and PHI, respectively).
Jin Zhao 0003, Yu Zhang 0027, Xiaofei Liao, Ligang He, Bingsheng He, Hai Jin 0001, Haikun Liu
SC4
2021 TurboDL: Improving the CNN Training on GPU With Fine-Grained Multi-Streaming Scheduling
abstract
Graphics Processing Units (GPUs) have evolved as powerful co-processors for the CNN training. Many new features have been introduced into GPUs such as concurrent kernel execution and hyper-Q technology. It is challenging to orchestrate concurrency for CNN (convolutional neural networks) training on GPUs since it may introduce synchronization overhead and poor resource utilization. Unlike previous research which mainly focuses on single layer or coarse-grained optimization, we introduce a critical-path based, asynchronous parallelization mechanism, and propose the optimization technique for the CNN training that takes into account global network architecture and GPU resource usage together. The proposed methods can effectively overlap the synchronization and the computation in different streams. As a result, the training process of CNN is accelerated. We have integrated our methods into Caffe. The experimental results show that the Caffe integrated with our methods can achieve 1.30X performance speedup on average compared with Caffe+cuDNN, and even higher performance speedup can be achieved for deeper, wider, and more complicated networks.
Hai Jin 0001, Xuanhua Shi, Ligang He, Bing Bing Zhou
IEEE Trans. Computers4
2021 SAFA: A Semi-Asynchronous Protocol for Fast Federated Learning With Low Overhead
abstract
Federated learning (FL) has attracted increasing attention as a promising approach to driving a vast number of end devices with artificial intelligence. However, it is very challenging to guarantee the efficiency of FL considering the unreliable nature of end devices while the cost of device-server communication cannot be neglected. In this article, we propose SAFA, a semi-asynchronous FL protocol, to address the problems in federated learning such as low round efficiency and poor convergence rate in extreme conditions (e.g., clients dropping offline frequently). We introduce novel designs in the steps of model distribution, client selection and global aggregation to mitigate the impacts of stragglers, crashes and model staleness in order to boost efficiency and improve the quality of the global model. We have conducted extensive experiments with typical machine learning tasks. The results demonstrate that the proposed protocol is effective in terms of shortening federated round duration, reducing local resource wastage, and improving the accuracy of the global model at an acceptable communication cost.
Wentai Wu, Ligang He, Weiwei Lin 0001, Rui Mao 0001, Carsten Maple, Stephen A. Jarvis
IEEE Trans. Computers2
2021 A Power Consumption Model for Cloud Servers Based on Elman Neural Network
abstract
Leveraging power consumption models in software systems can achieve easy deployment of low-cost, high-availability power monitoring in cloud datacenters that are usually large-scale, heterogeneous and frequently scaling up. However, traditional regression-based power consumption models generally have two drawbacks. First, their mathematical forms are usually fixed and determined a priori. This may cause unacceptable increase of error or over-fitting as the power signatures of cloud servers are usually uncertain. Second, the characteristic of workload dispatched to cloud servers is constantly changing while regression-based models can hardly generalize to a wide range of servers and workload types. As a novel solution, we in this paper propose a server power consumption model based on Elman Neural Network (PCM-ENN), aiming to allow accurate and flexible power estimation. PCM-ENN is an end-to-end black box model capable of learning the temporal relation between samples in a time series of power consumption. We trained and evaluated PCM-ENN on two power sequence datasets collected from heterogeneous hardware and operating systems running quasi-production benchmarks like CloudSuite. Experimental result shows that PCM-ENN generated accurate estimates on server power consumption with only small errors, outperforming widely-used linear regression model and NARX model in terms of accuracy.
Wentai Wu, Weiwei Lin 0001, Ligang He, Guangxin Wu, Ching-Hsien Hsu
IEEE Trans. Cloud Comput.3
2021 An Efficiency-Boosting Client Selection Scheme for Federated Learning With Fairness Guarantee
abstract
The issue of potential privacy leakage during centralized AI's model training has drawn intensive concern from the public. A Parallel and Distributed Computing (or PDC) scheme, termed Federated Learning (FL), has emerged as a new paradigm to cope with the privacy issue by allowing clients to perform model training locally, without the necessity to upload their personal sensitive data. In FL, the number of clients could be sufficiently large, but the bandwidth available for model distribution and re-upload is quite limited, making it sensible to only involve part of the volunteers to participate in the training process. The client selection policy is critical to an FL process in terms of training efficiency, the final model's quality as well as fairness. In this article, we will model the fairness guaranteed client selection as a Lyapunov optimization problem and then a C2MAB-based method is proposed for estimation of the model exchange time between each client and the server, based on which we design a fairness guaranteed algorithm termed RBCS-F for problem-solving. The regret of RBCS-F is strictly bounded by a finite constant, justifying its theoretical feasibility. Barring the theoretical results, more empirical data can be derived from our real training experiments on public datasets.
Tiansheng Huang, Weiwei Lin 0001, Wentai Wu, Ligang He, Keqin Li 0001, Albert Y. Zomaya
IEEE Trans. Parallel Distributed Syst.4
2021 Accelerating Federated Learning Over Reliability-Agnostic Clients in Mobile Edge Computing Systems
abstract
Mobile Edge Computing (MEC), which incorporates the Cloud, edge nodes, and end devices, has shown great potential in bringing data processing closer to the data sources. Meanwhile, Federated learning (FL) has emerged as a promising privacy-preserving approach to facilitating AI applications. However, it remains a big challenge to optimize the efficiency and effectiveness of FL when it is integrated with the MEC architecture. Moreover, the unreliable nature (e.g., stragglers and intermittent drop-out) of end devices significantly slows down the FL process and affects the global model's quality in such circumstances. In this article, a multi-layer federated learning protocol called HybridFL is designed for the MEC architecture. HybridFL adopts two levels (the edge level and the cloud level) of model aggregation enacting different aggregation strategies. Moreover, in order to mitigate stragglers and end device drop-out, we introduce regional slack factors into the stage of client selection performed at the edge nodes using a probabilistic approach without identifying or probing the state of end devices (whose reliability is agnostic). We demonstrate the effectiveness of our method in modulating the proportion of clients selected and present the convergence analysis for our protocol. We have conducted extensive experiments with machine learning tasks in different scales of MEC system. The results show that HybridFL improves the FL training process significantly in terms of shortening the federated round length, speeding up the global model's convergence (by up to 12×) and reducing end device energy consumption (by up to 58 percent).
Wentai Wu, Ligang He, Weiwei Lin 0001, Rui Mao 0001
IEEE Trans. Parallel Distributed Syst.2
2021 Feluca: A Two-Stage Graph Coloring Algorithm With Color-Centric Paradigm on GPU
abstract
There are great challenges in performing graph coloring on GPU in general. First, the long-tail problem exists in the recursion algorithm because the conflict (i.e., different threads assign the adjacent nodes to the same color) becomes more likely to occur as the number of iterations increases. Second, it is hard to parallelize the sequential spread algorithm because the color allocation depends on the adjoining iteration. Third, the atomic operation is widely used on GPU to maintain the color list, which can greatly reduce the efficiency of GPU threads. In this article, we propose a two-stage high-performance graph coloring algorithm, called Feluca, aiming to address the above challenges. Feluca combines the recursion-based method with the sequential spread-based method. In the first stage, Feluca uses a recursive routine to color a majority of vertices in the graph. Then, it switches to the sequential spread method to color the remaining vertices in order to avoid the conflicts of the recursive algorithm. Moreover, the following techniques are proposed to further improve the graph coloring performance. i) A new method is proposed to eliminate the cycles in the graph; ii) a top-down scheme is developed to avoid the atomic operation originally required for color selection; and iii) a novel color-centric coloring paradigm is designed to improve the degree of parallelism for the sequential spread part. All these newly developed techniques, together with further GPU-specific optimizations such as coalesced memory access, comprise an efficient parallel graph coloring solution in Feluca. We have conducted extensive experiments on NVIDIA GPU. The results show that Feluca can achieve 1.19 - 8.39× speedup over the state-of-the-art algorithms.
Zhigao Zheng 0001, Xuanhua Shi, Ligang He, Hai Jin 0001, Shuo Wei, Hulin Dai
IEEE Trans. Parallel Distributed Syst.3
2020 Developing a Loss Prediction-based Asynchronous Stochastic Gradient Descent Algorithm for Distributed Training of Deep Neural Networks
abstract
Training Deep Neural Network is a computation-intensive and time-consuming task. Asynchronous Stochastic Gradient Descent (ASGD) is an effective solution to accelerate the training process since it enables the network to be trained in a distributed fashion, but with a main issue of the delayed gradient update. A recent notable work called DC-ASGD improves the performance of ASGD by compensating the delay using a cheap approximation of the Hessian matrix. DC-ASGD works well with a short delay; however, the performance drops considerably with an increasing delay between the workers and the server. In real-life large-scale distributed training, such gradient delay experienced by the worker is usually high and volatile. In this paper, we propose a novel algorithm called LC-ASGD to compensate for the delay, basing on Loss Prediction. It effectively extends the tolerable delay duration for the compensation mechanism. Specifically, LC-ASGD utilizes additional models that reside in the parameter server and predict the loss to compensate for the delay, basing on historical losses collected from each worker. The algorithm is evaluated on the popular networks and benchmark datasets. The experimental results show that our LC-ASGD significantly improves over existing methods, especially when the networks are trained with a large number of workers.
Ligang He, Shenyuan Ren, Rui Mao 0001
ICPP2
2020 A parameter-level parallel optimization algorithm for large-scale spatio-temporal data mining
Xuanhua Shi, Ligang He, Dongxiao Yu, Hai Jin 0001, Chen Yu 0003, Hulin Dai, Zezhao Feng
Distributed Parallel Databases3
2020 WolfGraph: The edge-centric graph processing on GPU
Huanzhou Zhu, Ligang He, Matthew Leeke, Rui Mao 0001
Future Gener. Comput. Syst.2
2020 A reformed task scheduling algorithm for heterogeneous distributed systems with energy consumption constraints
Yikun Hu 0001, Jinghong Li, Ligang He
Neural Comput. Appl.3
2020 Optimizing the SSD Burst Buffer by Traffic Detection
abstract
Currently, HPC storage systems still use hard disk drive (HDD) as their dominant storage device. Solid state drive (SSD) is widely deployed as the buffer to HDDs. Burst buffer has also been proposed to manage the SSD buffering of bursty write requests. Although burst buffer can improve I/O performance in many cases, we find that it has some limitations such as requiring large SSD capacity and harmonious overlapping between computation phase and data flushing phase. In this article, we propose a scheme, called SSDUP+. 1 SSDUP+ aims to improve the burst buffer by addressing the above limitations. First, to reduce the demand for the SSD capacity, we develop a novel method to detect and quantify the data randomness in the write traffic. Further, an adaptive algorithm is proposed to classify the random writes dynamically. By doing so, much less SSD capacity is required to achieve the similar performance as other burst buffer schemes. Next, to overcome the difficulty of perfectly overlapping the computation phase and the flushing phase, we propose a pipeline mechanism for the SSD buffer, in which data buffering and flushing are performed in pipeline. In addition, to improve the I/O throughput, we adopt a traffic-aware flushing strategy to reduce the I/O interference in HDD. Finally, to further improve the performance of buffering random writes in SSD, SSDUP+ transforms the random writes to sequential writes in SSD by storing the data with a log structure. Further, SSDUP+ uses the AVL tree structure to store the sequence information of the data. We have implemented a prototype of SSDUP+ based on OrangeFS and conducted extensive experiments. The experimental results show that our proposed SSDUP+ can save an average of 50% SSD space while delivering almost the same performance as other common burst buffer schemes. In addition, SSDUP+ can save about 20% SSD space compared with the previous version of this work, SSDUP, while achieving 20–30% higher I/O throughput than SSDUP.
Xuanhua Shi, Wei Liu 0004, Ligang He, Hai Jin 0001, Yong Chen 0001
ACM Trans. Archit. Code Optim.3
2019 Developing the Parallelization Methods for Finding the All-Pairs Shortest Paths in Distributed Memory Architecture
abstract
The All-Pairs Shortest Paths (APSP) is a fundamental graph problem aiming to find the shortest path between any two nodes in a graph. In this paper, a new method is presented to solve the APSP problem for big graphs on distributed systems. In this method, a graph is partitioned judiciously and then processed in parallel. In particular, the graph is first pre-processed to prepare the partition in the computation stages. After the graph is partitioned into smaller sub-graphs, a traditional shortest path algorithm, such as the Floyd-Warshall algorithm or the Dijkstra's algorithm, can be used to find the APSP in each sub-graph. Finally, through the common nodes between the sub-graphs, the local results in each sub-graph are combined to establish the APSP for the entire graph. Our method is implemented with MPI. Two different communication patterns among partitions (and processes) are proposed to achieve the parallelization and the combination of the local results. We then conducted extensive experiments on a high performance cluster. The experimental results show that comparing with the existing solution, our method is able to accelerate the solving of the APSP problem significantly.
Mohammed Alghamdi, Ligang He, Yujue Zhou
IPCCC2
2019 GraphM: an efficient storage system for high throughput of concurrent graph processing
abstract
With the rapidly growing demand of graph processing in the real world, a large number of iterative graph processing jobs run concurrently on the same underlying graph. However, the storage engines of existing graph processing frameworks are mainly designed for running an individual job. Our studies show that they are inefficient when running concurrent jobs due to the redundant data storage and access overhead. To cope with this issue, we develop an efficient storage system, called GraphM. It can be integrated into the existing graph processing systems to efficiently support concurrent iterative graph processing jobs for higher throughput by fully exploiting the similarities of the data accesses between these concurrent jobs. GraphM regularizes the traversing order of the graph partitions for concurrent graph processing jobs by streaming the partitions into the main memory and the Last-Level Cache (LLC) in a common order, and then processes the related jobs concurrently in a novel fine-grained synchronization. In this way, the concurrent jobs share the same graph structure data in the LLC/memory and also the data accesses to the graph, so as to amortize the storage consumption and the data access overhead. To demonstrate the efficiency of GraphM, we plug it into state-of-the-art graph processing systems, including GridGraph, GraphChi, PowerGraph, and Chaos. Experiments results show that GraphM improves the throughput by 1.73~13 times.
Jin Zhao 0003, Yu Zhang 0027, Xiaofei Liao, Ligang He, Bingsheng He, Hai Jin 0001, Haikun Liu
SC4
2019 Software-defined QoS for I/O in exascale computing
Yusheng Hua, Xuanhua Shi, Hai Jin 0001, Wei Liu 0004, Yong Chen 0001, Ligang He
CCF Trans. High Perform. Comput.7
2019 CGraph: A Distributed Storage and Processing System for Concurrent Iterative Graph Analysis Jobs
abstract
Distributed graph processing platforms usually need to handle massive Concurrent iterative Graph Processing (CGP) jobs for different purposes. However, existing distributed systems face high ratio of data access cost to computation for the CGP jobs, which incurs low throughput. We observed that there are strong spatial and temporal correlations among the data accesses issued by different CGP jobs, because these concurrently running jobs usually need to repeatedly traverse the shared graph structure for the iterative processing of each vertex. Based on this observation, this article proposes a distributed storage and processing system CGraph for the CGP jobs to efficiently handle the underlying static/evolving graph for high throughput. It uses a data-centric load-trigger-pushing model, together with several optimizations, to enable the CGP jobs to efficiently share the graph structure data in the cache/memory and their accesses by fully exploiting such correlations, where the graph structure data is decoupled from the vertex state associated with each job. It can deliver much higher throughput for the CGP jobs by effectively reducing their average ratio of data access cost to computation. Experimental results show that CGraph improves the throughput of the CGP jobs by up to 3.47× in comparison with existing solutions on distributed platforms.
Yu Zhang 0027, Jin Zhao 0003, Xiaofei Liao, Hai Jin 0001, Lin Gu 0002, Haikun Liu, Bingsheng He, Ligang He
ACM Trans. Storage8
2018 vPlacer: A Co-scheduler for Optimizing the Performance of Parallel Jobs in Xen
Ligang He, Shenyuan Ren, Rui Mao 0001
ICA3PP (1)2
2018 Scheduling DAG Applications for Time Sharing Systems
Shenyuan Ren, Ligang He, Chao Chen 0011, Zhuoer Gu
ICA3PP (2)2
2018 vGrouper: Optimizing the Performance of Parallel Jobs in Xen by Increasing Synchronous Execution of Virtual Machines
Ligang He, Shenyuan Ren, Yuhua Cui
NPC2
2018 Data Fine-Pruning: A Simple Way to Accelerate Neural Network Training
Ligang He, Shenyuan Ren, Rui Mao 0001
NPC2
2018 CGraph: A Correlations-aware Approach for Efficient Concurrent Iterative Graph Processing
Yu Zhang 0027, Xiaofei Liao, Hai Jin 0001, Lin Gu 0002, Ligang He, Bingsheng He, Haikun Liu
USENIX ATC5
2018 Concurrent hash tables on multicore machines: Comparison, evaluation and implications
Zhiwen Chen 0006, Xin He 0054, Jianhua Sun 0002, Hao Chen 0002, Ligang He
Future Gener. Comput. Syst.5
2018 Modelling and developing conflict-aware scheduling on large-scale data centres
Chao Chen 0011, Ligang He, Bo Gao 0001, Jiadong Ren, Zhangjie Fu 0001, Songling Fu, Yongjian Hu, Chang-Tsun Li
Future Gener. Comput. Syst.3
2018 Energy-efficient hadoop for big data analytics and computing: A systematic review and research insights
Wentai Wu, Weiwei Lin 0001, Ching-Hsien Hsu, Ligang He
Future Gener. Comput. Syst.4
2018 Deca: A Garbage Collection Optimizer for In-Memory Data Processing
abstract
In-memory caching of intermediate data and active combining of data in shuffle buffers have been shown to be very effective in minimizing the recomputation and I/O cost in big data processing systems such as Spark and Flink. However, it has also been widely reported that these techniques would create a large amount of long-living data objects in the heap. These generated objects may quickly saturate the garbage collector, especially when handling a large dataset, and hence, limit the scalability of the system. To eliminate this problem, we propose a lifetime-based memory management framework, which, by automatically analyzing the user-defined functions and data types, obtains the expected lifetime of the data objects and then allocates and releases memory space accordingly to minimize the garbage collection overhead. In particular, we present Deca,1 a concrete implementation of our proposal on top of Spark, which transparently decomposes and groups objects with similar lifetimes into byte arrays and releases their space altogether when their lifetimes come to an end. When systems are processing very large data, Deca also provides field-oriented memory pages to ensure high compression efficiency. Extensive experimental studies using both synthetic and real datasets show that, in comparing to Spark, Deca is able to (1) reduce the garbage collection time by up to 99.9%, (2) reduce the memory consumption by up to 46.6% and the storage space by 23.4%, (3) achieve 1.2× to 22.7× speedup in terms of execution time in cases without data spilling and 16× to 41.6× speedup in cases with data spilling, and (4) provide similar performance compared to domain-specific systems.
Xuanhua Shi, Zhixiang Ke, Yongluan Zhou, Hai Jin 0001, Lu Lu 0006, Ligang He
ACM Trans. Comput. Syst.7
2017 MURS: Mitigating Memory Pressure in Service-Oriented Data Processing System
abstract
Although a data processing system often works as a batch processing system, many enterprises deploy such a system as a service, which we call the service-oriented data processing system. It has been shown that in-memory data processing systems suffer from serious memory pressure. The situation becomes even worse for the service-oriented data processing systems due to various reasons. For example, in a service-oriented system, multiple submitted tasks are launched at the same time and executed in the same context in the resources, compared with the batch processing mode where the tasks are processed one by one. Therefore, the memory pressure will affect all submitted tasks, including the tasks that only incur the light memory pressure when they are run alone. In this paper, we find that the reason why memory pressure arises is because the running tasks produce massive long-living data objects in the limited memory space. Our studies further reveal that the long-living data objects are generated by the API functions that are invoked by the in-memory processing frameworks. Based on these findings, we propose a method to classify the API functions based on the memory usage rate. Further, we design a scheduler called MURS to mitigate the memory pressure. We implement MURS in Spark and conduct the experiments to evaluate the performance of MURS. The results show that when comparing to Spark, MURS can 1) decrease the execution time of the submitted jobs by up to 65.8%, 2) mitigate the memory pressure in the server by decreasing the garbage collection time by up to 81%, and 3) reduce the data spilling, and hence disk I/O, by approximately 90%.
Xuanhua Shi, Ligang He, Hai Jin 0001, Zhixiang Ke, Song Wu 0001
ICWS3
2017 Developing power-aware scheduling mechanisms for computing systems virtualized by Xen
abstract
Summary Cloud computing emerges as one of the most important technologies for interconnecting people and building the so‐called Internet of People (IoP). In such a cloud‐based IoP, the virtualization technique provides the key supporting environments for running the IoP jobs such as performing data analysis and mining personal information. Nowadays, energy consumption in such a system is a critical metric to measure the sustainability and eco‐friendliness of the system. This paper develops three power‐aware scheduling strategies in virtualized systems managed by Xen, which is a popular virtualization technique. These three strategies are the Least performance Loss Scheduling strategy, the No performance Loss Scheduling strategy, and the Best Frequency Match scheduling strategy. These power‐aware strategies are developed by identifying the limitation of Xen in scaling the CPU frequency and aim to reduce the energy waste without sacrificing the jobs running performance in the computing systems virtualized by Xen. Least performance Loss Scheduling works by re‐arranging the execution order of the virtual machines (VMs). No performance Loss Scheduling works by setting a proper initial CPU frequency for running the VMs. Best Frequency Match reduces energy waste and performance loss by allowing the VMs to jump the queue so that the VM that is put into execution best matches the current CPU frequency. Scheduling for both single core and multicore processors is considered in this paper. The evaluation experiments have been conducted, and the results show that compared with the original scheduling strategy in Xen, the developed power‐aware scheduling algorithm is able to reduce energy consumption without reducing the performance for the jobs running in Xen. Copyright © 2016 John Wiley & Sons, Ltd.
Shenyuan Ren, Ligang He, Huanzhou Zhu, Zhuoer Gu, Jiandong Shang
Concurr. Comput. Pract. Exp.2
2017 Performance analysis and optimization for workflow authorization
abstract
Cloud download service, as a new application which downloads the requested content offline and reserves it in cloud storage until users retrieve it, has recently become a trend attracting millions of users in China. In the face of the dilemma between the growth of download requests and the limitation of storage resource, the cloud servers have to design an efficient resource allocation scheme to enhance the utilization of storage as well as to satisfy users' needs like a short download time. When a user's churn behavior is considered as a Markov chain process, it is found that a proper allocation of download speed can optimize the storage resource utilization. Accordingly, two dynamic resource allocation schemes including a speed switching (SS) scheme and a speed increasing (SI) scheme are proposed. Both theoretical analysis and simulation results prove that our schemes can effectively reduce the consumption of storage resource and keep the download time short enough for a good user experience.
Ligang He, Nadeem Chaudhary, Songling Fu, Hao Chen 0002, Jianhua Sun 0002, Kenli Li 0001, Zhangjie Fu 0001
Future Gener. Comput. Syst.2
2017 Robot Cloud: Bridging the power of robotics and cloud computing
Zhihui Du, Ligang He, Yinong Chen 0004, Tongzhou Wang 0002
Future Gener. Comput. Syst.2
2017 Multi-resource scheduling and power simulation for cloud computing
Weiwei Lin 0001, Siyao Xu, Ligang He, Jin Li 0002
Inf. Sci.3
2016 Developing the Cloud-integrated data replication framework in decentralized online social networks
Songling Fu, Ligang He, Xiangke Liao, Chenlin Huang
J. Comput. Syst. Sci.2
2016 Lifetime-Based Memory Management for Distributed Data Processing Systems
abstract
In-memory caching of intermediate data and eager combining of data in shuffle buffers have been shown to be very effective in minimizing the re-computation and I/O cost in distributed data processing systems like Spark and Flink. However, it has also been widely reported that these techniques would create a large amount of long-living data objects in the heap, which may quickly saturate the garbage collector, especially when handling a large dataset, and hence would limit the scalability of the system. To eliminate this problem, we propose a lifetime-based memory management framework, which, by automatically analyzing the user-defined functions and data types, obtains the expected lifetime of the data objects, and then allocates and releases memory space accordingly to minimize the garbage collection overhead. In particular, we present Deca, a concrete implementation of our proposal on top of Spark, which transparently decomposes and groups objects with similar lifetimes into byte arrays and releases their space altogether when their lifetimes come to an end. An extensive experimental study using both synthetic and real datasets shows that, in comparing to Spark, Deca is able to 1) reduce the garbage collection time by up to 99.9%, 2) to achieve up to 22.7x speed up in terms of execution time in cases without data spilling and 41.6x speedup in cases with data spilling, and 3) to consume up to 46.6% less memory.
Lu Lu 0006, Xuanhua Shi, Yongluan Zhou, Hai Jin 0001, Cheng Pei, Ligang He, Yuanzhen Geng
Proc. VLDB Endow.7
2016 Developing Graph-Based Co-Scheduling Algorithms on Multicore Computers
abstract
It is common that multiple cores reside on the same chip and share the on-chip cache. As a result, resource sharing can cause performance degradation of co-running jobs. Job co-scheduling is a technique that can effectively alleviate this contention and many co-schedulers have been reported in related literature. Most solutions however do not aim to find the optimal co-scheduling solution. Being able to determine the optimal solution is critical for evaluating co-scheduling systems. Moreover, most co-schedulers only consider serial jobs, and there often exist both parallel and serial jobs in real-world systems. In this paper a graph-based method is developed to find the optimal co-scheduling solution for serial jobs; the method is then extended to incorporate parallel jobs, including multi-process, and multi-threaded parallel jobs. A number of optimization measures are also developed to accelerate the solving process. Moreover, a flexible approximation technique is proposed to strike a balance between the solving speed and the solution quality. Extensive experiments are conducted to evaluate the effectiveness of the proposed co-scheduling algorithms. The results show that the proposed algorithms can find the optimal co-scheduling solution for both serial and parallel jobs. The proposed approximation technique is also shown to be flexible in the sense that we can control the solving speed by setting the requirement for the solution quality.
Ligang He, Huanzhou Zhu, Stephen A. Jarvis
IEEE Trans. Parallel Distributed Syst.1
2016 Redundant Network Traffic Elimination with GPU Accelerated Rabin Fingerprinting
abstract
Recently, redundant network traffic elimination has attracted a lot of attention from both the academia and the industry. A core challenge and enabling technique in implementing redundancy elimination is to perform content-based chunking, which typically involves the computationally heavy Rabin fingerprinting algorithm. In this paper, we propose a GPU-based implementation of Rabin fingerprinting to address this issue. To maximize performance gains, a diverse set of optimization strategies, such as efficient buffer management, GPU memory hierarchy optimization, and balanced load distribution, is proposed by either exploiting the intrinsic hardware features or addressing domain-specific challenges. Extensive evaluations on both the overall and microscopic performance reveal the effectiveness of the GPU-accelerated Rabin fingerprinting algorithm, and we can achieve up to 40 Gpbs throughput on a GTX 780 card. The throughput shows 1.87× speedup against the state-of-the-art using comparable hardware. In addition, although some optimization designs are specific for the problem, techniques proposed in this work including the indexed compact buffer scheme and approximate sorting would also be beneficial and applicable to other network applications leveraging GPU acceleration.
Jianhua Sun 0002, Hao Chen 0002, Ligang He, Huailiang Tan
IEEE Trans. Parallel Distributed Syst.3
2015 GPSA: A Graph Processing System with Actors
abstract
Due to the increasing need to process the fast growing graph-structured data (e.g. Social networks and Web graphs), designing high performance graph processing systems becomes one of the most urgent problems facing systems researchers. In this paper, we introduce GPSA, a single-machine graph processing system based on an actor computation model inspired by the Bulk Synchronous Parallel(BSP) computation model. GPSA takes advantage of actors to improve the concurrency on a single machine with limited resource. GPSA improves the conventional BSP computation model to fit in the actor programming paradigm by decoupling the message dispatching from the computation. Furthermore, we exploit memory mapping to avoid explicit data management to improve I/O performance. Experimental evaluation shows that our system outperforms existing systems by 2x-6x in processing large-scale graphs on a single system.
Jianhua Sun 0002, Dongwei Zhou, Hao Chen 0002, Zhiwen Chen 0006, Ligang He
ICPP7
2015 Modelling and Developing Co-scheduling Strategies on Multicore Processors
abstract
On-chip cache is often shared between processes that run concurrently on different cores of the same processor. Resource contention of this type causes performance degradation to the co-running processes. Contention-aware co-scheduling refers to the class of scheduling techniques to reduce the performance degradation. Most existing contention-aware co-schedulers only consider serial jobs. However, there often exist both parallel and serial jobs in computing systems. In this paper, the problem of co-scheduling a mix of serial and parallel jobs is modelled as an Integer Programming (IP) problem. Then the existing IP solver can be used to find the optimal co-scheduling solution that minimizes the performance degradation. However, we find that the IP-based method incurs high time overhead and can only be used to solve small-scale problems. Therefore, a graph-based method is also proposed in this paper to tackle this problem. We construct a co-scheduling graph to represent the co-scheduling problem and model the problem of finding the optimal co-scheduling solution as the problem of finding the shortest valid path in the co-scheduling graph. A heuristic A*-search algorithm (HA*) is then developed to find the near-optimal solutions efficiently. The extensive experiments have been conducted to verify the effectiveness and efficiency of the proposed methods. The experimental results show that compared with the IP-based method, HA* is able to find the near-optimal solutions with much less time.
Huanzhou Zhu, Ligang He, Bo Gao 0001, Kenli Li 0001, Jianhua Sun 0002, Hao Chen 0002, Keqin Li 0001
ICPP2
2015 Modelling and Optimizing Bandwidth Provision for Interacting Cloud Services
Chao Chen 0011, Ligang He, Bo Gao 0001, Kenli Li 0001, Keqin Li 0001
ICSOC2
2015 Modelling the Bandwidth Allocation Problem in Mobile Service-Oriented Networks
abstract
When the services requested by mobile application workflows are distributed over a network of mobile smart devices, the question arises as to which service should be allocated with how much bandwidth and when in order to satisfy service demands? Furthermore, the mobility of smart mobile devices brings forward the challenge to determine how changes in mobile network conditions affect the bandwidth requirements of interacting services. In this paper, we construct a Network I-O model to describe the bandwidth dependencies in mobile service-oriented networks incorporating and extend on the principles of the Leontief Input-Output model in economics. Various factors such as bandwidth and service demand are accounted for in the model. The network I-O model lays the foundation for future objective developments in ubiquitous mobile computing scenarios. Results from simulation studies are presented to demonstrate the effectiveness of the proposed methods.
Bo Gao 0001, Ligang He, Chao Chen 0011
MSWiM2
2015 Performance Optimization for Managing Massive Numbers of Small Files in Distributed File Systems
abstract
The processing of massive numbers of small files is a challenge in the design of distributed file systems. Currently, the combined-block-storage approach is prevalent. However, the approach employs the traditional file systems such as ExtFS and may cause inefficiency when accessing small files randomly located in the disk. This paper focuses on optimizing the performance of data servers in accessing massive numbers of small files. We present a Flat Lightweight File System (iFlatLFS) to manage small files, which is based on a simple metadata scheme and a flat storage architecture. iFlatLFS is designed to substitute the traditional file system on data servers and can be deployed underneath distributed file systems that store massive numbers of small files. iFlatLFS can greatly simplify the original data access procedure. The new metadata proposed in this paper occupies only a fraction of the metadata size based on traditional file systems. We have implemented iFlatLFS in CentOS 5.5 and integrated it into an open source Distributed File System (DFS), called Taobao FileSystem (TFS), which is developed by a top B2C service provider, Alibaba, in China and is managing over 28.6 billion small photos. We have conducted extensive experiments to verify the performance of iFlatLFS. The results show that when the file size ranges from 1 to 64 KB, iFlatLFS is faster than Ext4 by 48 and 54 percent on average for random read and write in the DFS environment, respectively. Moreover, after iFlatLFS is integrated into TFS, iFlatLFS-based TFS is faster than the existing Ext4-based TFS by 45 and 49 percent on average for random read access and hybrid access (the mix of read and write accesses), respectively.
Songling Fu, Ligang He, Chenlin Huang, Xiangke Liao, Kenli Li 0001
IEEE Trans. Parallel Distributed Syst.2
2015 Mammoth: Gearing Hadoop Towards Memory-Intensive MapReduce Applications
abstract
The MapReduce platform has been widely used for large-scale data processing and analysis recently. It works well if the hardware of a cluster is well configured. However, our survey has indicated that common hardware configurations in small- and medium-size enterprises may not be suitable for such tasks. This situation is more challenging for memory-constrained systems, in which the memory is a bottleneck resource compared with the CPU power and thus does not meet the needs of large-scale data processing. The traditional high performance computing (HPC) system is an example of the memory-constrained system according to our survey. In this paper, we have developed Mammoth, a new MapReduce system, which aims to improve MapReduce performance using global memory management. In Mammoth, we design a novel rule-based heuristic to prioritize memory allocation and revocation among execution units (mapper, shuffler, reducer, etc.), to maximize the holistic benefits of the Map/Reduce job when scheduling each memory unit. We have also developed a multi-threaded execution engine, which is based on Hadoop but runs in a single JVM on a node. In the execution engine, we have implemented the algorithm of memory scheduling to realize global memory management, based on which we further developed the techniques such as sequential disk accessing, multi-cache and shuffling from memory, and solved the problem of full garbage collection in the JVM. We have conducted extensive experiments to compare Mammoth against the native Hadoop platform. The results show that the Mammoth system can reduce the job execution time by more than 40 percent in typical cases, without requiring any modifications of the Hadoop programs. When a system is short of memory, Mammoth can improve the performance by up to 5.19 times, as observed for I/O intensive applications, such as PageRank. We also compared Mammoth with Spark. Although Spark can achieve better performance than Mammoth for interactive and iterative applications when the memory is sufficient, our experimental results show that for batch processing applications, Mammoth can adapt better to various memory environments and outperform Spark when the memory is insufficient, and can obtain similar performance as Spark when the memory is sufficient. Given the growing importance of supporting large-scale data processing and analysis and the proven success of the MapReduce platform, the Mammoth system can have a promising potential and impact.
Xuanhua Shi, Ligang He, Lu Lu 0006, Hai Jin 0001, Yong Chen 0001, Song Wu 0001
IEEE Trans. Parallel Distributed Syst.3
2015 A Hybrid Chemical Reaction Optimization Scheme for Task Scheduling on Heterogeneous Computing Systems
abstract
Scheduling for directed acyclic graph (DAG) tasks with the objective of minimizing makespan has become an important problem in a variety of applications on heterogeneous computing platforms, which involves making decisions about the execution order of tasks and task-to-processor mapping. Recently, the chemical reaction optimization (CRO) method has proved to be very effective in many fields. In this paper, an improved hybrid version of the CRO method called HCRO (hybrid CRO) is developed for solving the DAG-based task scheduling problem. In HCRO, the CRO method is integrated with the novel heuristic approaches, and a new selection strategy is proposed. More specifically, the following contributions are made in this paper. (1) A Gaussian random walk approach is proposed to search for optimal local candidate solutions. (2) A left or right rotating shift method based on the theory of maximum Hamming distance is used to guarantee that our HCRO algorithm can escape from local optima. (3) A novel selection strategy based on the normal distribution and a pseudo-random shuffle approach are developed to keep the molecular diversity. Moreover, an exclusive-OR (XOR) operator between two strings is introduced to reduce the chance of cloning before new molecules are generated. Both simulation and real-life experiments have been conducted in this paper to verify the effectiveness of HCRO. The results show that the HCRO algorithm schedules the DAG tasks much better than the existing algorithms in terms of makespan and speed of convergence.
Yuming Xu, Kenli Li 0001, Ligang He, Longxin Zhang, Keqin Li 0001
IEEE Trans. Parallel Distributed Syst.3
2014 Modelling and Predicting the Data Availability in Decentralized Online Social Networks
abstract
Maintaining data availability is one of the biggest challenges in Decentralized Online Social Networks (DOSN). In the existing work of improving data availability in DOSN, it is often assumed that the friends of a user are always capable of contributing sufficient storage capacity to store all the data published by the user. However, this assumption is not always true for today's Online Social Networks (OSNs) for the following reasons. On one hand, the increasingly more data are being generated on the OSNs nowadays. On the other hand, current users often use the smart mobile devices to access the OSNs. These two factors cause the shortage of the storage capacity in DOSN, where the published data are supposed to be stored within a friend circle. The limitation of the storage capacity may jeopardize the data availability. Therefore, it is desired to know the relation between the storage capacity contributed by the OSN users and the level of data availability that the OSN can achieve. This paper addresses this issue. In this paper, the data availability model over storage capacity is established. Further, a novel method is proposed to predict the data availability on the fly. Extensive simulation experiments have been conducted to evaluate the effectiveness of the data availability model and the on-the-fly prediction. The data availability model can be used by the OSN designers to determine the storage capacity for the published data in order to achieve the desired data availability. The on-the-fly prediction method can help the data replication and storage policies make judicious decisions at runtime.
Songling Fu, Ligang He, Xiangke Liao, Chenlin Huang, Kenli Li 0001, Bo Gao 0001
ICWS2
2014 Optimizing Job Scheduling on Multicore Computers
abstract
It is common nowadays that multiple cores reside on the same chip and share the on-chip cache. Resource sharing may cause performance degradation of the co-running jobs. Job co-scheduling is a technique that can effectively alleviate the contention. Many co-schedulers have been developed in the literature, but most of them do not aim to find the optimal co-scheduling solution. Being able to determine the optimal solution is critical for evaluating co-scheduling systems. Moreover, most co-schedulers only consider serial jobs. However, there often exist both parallel and serial jobs in some situations. This paper aims to tackle these issues. In this paper, a graph-based method is developed to find the optimal co-scheduling solution for serial jobs, and then the method is extended to incorporate parallel jobs. The extensive experiments have been conducted to evaluate the effectiveness and efficiency of the proposed co-scheduling algorithms. The results show that the proposed algorithms can find the optimal co-scheduling solution for both serial and parallel jobs.
Huanzhou Zhu, Ligang He, Stephen A. Jarvis
MASCOTS2
2014 Developing security-aware resource management strategies for workflows
Ligang He, Nadeem Chaudhary, Stephen A. Jarvis
Future Gener. Comput. Syst.1
2014 Developing resource consolidation frameworks for moldable virtual machines in clouds
Ligang He, Deqing Zou, Chao Chen 0011, Hai Jin 0001, Stephen A. Jarvis
Future Gener. Comput. Syst.1
2014 BAG: Managing GPU as Buffer Cache in Operating Systems
abstract
This paper presents the design, implementation and evaluation of BAG, a system that manages GPU as the buffer cache in operating systems. Unlike previous uses of GPUs, which have focused on the computational capabilities of GPUs, BAG is designed to explore a new dimension in managing GPUs in heterogeneous systems where the GPU memory is an exploitable but always ignored resource. With the carefully designed data structures and algorithms, such as concurrent hashtable, log-structured data store for the management of GPU memory, and highly-parallel GPU kernels for garbage collection, BAG achieves good performance under various workloads. In addition, leveraging the existing abstraction of the operating system not only makes the implementation of BAG non-intrusive, but also facilitates the system deployment.
Hao Chen 0002, Jianhua Sun 0002, Ligang He, Kenli Li 0001, Huailiang Tan
IEEE Trans. Parallel Distributed Syst.3
2013 Developing communication-aware service placement frameworks in the Cloud economy
abstract
In a Cloud system, a number of services are often deployed with each service being hosted by a collection of Virtual Machines (VM). The services may interact with each other and the interaction patterns may be dynamic, varying according to the system information at runtime. These impose a challenge in determining the amount of resources required to deliver a desired level of QoS for each service. In this paper, we present a method to determine the sufficient number of VMs for the interacting Cloud services. The proposed method borrows the ideas from the Leontief Open Production Model in economy. Further, this paper develops a communication-aware strategy to place the VMs to Physical Machines (PM), aiming to minimize the communication costs incurred by the service interactions. The developed communication-aware placement strategy is formalized in a way that it does not need to the specific communication pattern between individual VMs. A genetic algorithm is developed to find a VM-to-PM placement with low communication costs. Simulation experiments have been conducted to evaluate the performance of the developed communication-aware placement framework. The results show that compared with the placement framework aiming to use the minimal number of PMs to host VMs, the proposed communication-aware framework is able to reduce the communication cost significantly with only a very little increase in the PM usage.
Chao Chen 0011, Ligang He, Hao Chen 0002, Jianhua Sun 0002, Bo Gao 0001, Stephen A. Jarvis
CLUSTER2
2013 Analyzing the performance impact of authorization constraints and optimizing the authorization methods for workflows
abstract
Many workflow management systems have been developed to enhance the performance of workflow executions. The authorization policies deployed in the system may restrict the task executions. The common authorization constraints include role constraints, Separation of Duty (SoD), Binding of Duty (BoD) and temporal constraints. This paper presents the methods to check the feasibility of these constraints, and also determines the time durations when the temporal constraints will not impose negative impact on performance. Further, this paper presents an optimal authorization method, which is optimal in the sense that it can minimize a workflow's delay caused by the temporal constraints. Simulation experiments have been conducted to verify the effectiveness of the proposed authorization method. The experimental results show that comparing with the intuitive authorization method, the optimal authorization method can reduce the delay caused by the authorization constraints and consequently reduce the workflows' response time.
Nadeem Chaudhary, Ligang He
HiPC2
2013 iFlatLFS: Performance optimization for accessing massive small files
abstract
The processing of massive small files is a challenge in the design of distributed file systems. Currently, the combined-block-storage approach is prevalent. However, the approach employs traditional file systems like ExtFS and may cause inefficiency for random access to small files. This paper focuses on optimizing the performance of data servers in accessing massive small files. We present a Flat Lightweight File System (iFlatLFS) to manage small files, which is based on a simple metadata scheme and a flat storage architecture. iFlatLFS aims to substitute the traditional file system on data servers that are mainly used to store small files, and it can greatly simplify the original data access procedure. The new metadata proposed in this paper occupies only a fraction of the original metadata size based on traditional file systems. We have implemented iFlatLFS in CentOS 5.5 and integrated it into an open source Distributed File System (DFS), called Taobao FileSystem (TFS), which is developed by a top B2C service provider, Alibaba, in China and is managing over 28.6 billion small photos. We have conducted extensive experiments to verify the performance of iFlatLFS. The results show that when the file size ranges from 1KB to 64KB, iFlatLFS is faster than Ext4 by 48% and 54% on average for random read and write in the DFS environment, respectively. Moreover, after iFlatLFS is integrated into TFS, iFlatLFS-based TFS is faster than the existing Ext4-based TFS by 45% and 49% on average for random read access and hybrid access (the mix of read and write accesses), respectively.
Songling Fu, Chenlin Huang, Ligang He, Nadeem Chaudhary, Xiangke Liao, Shazhou Yang, Bao Li 0002
HiPC3
2013 Modelling Energy-Aware Task Allocation in Mobile Workflows
Bo Gao 0001, Ligang He
MobiQuitous2
2013 Towards Automated Memory Model Generation Via Event Tracing
abstract
The importance of memory performance and capacity is a growing concern for high performance computing laboratories around the world. It has long been recognized that improvements in processor speed exceed the rate of improvement in dynamic random access memory speed and, as a result, memory access times can be the limiting factor in high performance scientific codes. The use of multi-core processors exacerbates this problem with the rapid growth in the number of cores not being matched by similar improvements in memory capacity, increasing the likelihood of memory contention. In this paper, we present WMTools, a lightweight memory tracing tool and analysis framework for parallel codes, which is able to identify peak memory usage and also analyse per-function memory use over time. An evaluation of WMTools, in terms of its effectiveness and also its overheads, is performed using nine established scientific applications/benchmark codes representing a variety of programming languages and scientific domains. We also show how WMTools can be used to automatically generate a parameterized memory model for one of these applications, a two-dimensional non-linear magnetohydrodynamics application, Lare2D. Through the memory model we are able to identify an unexpected growth term which becomes dominant at scale. With a refined model we are able to predict memory consumption with under 7% error.
Oliver Perks, D. A. Beckingsale, Simon D. Hammond, I. Miller, J. A. Herdman, A. Vadgama, Abhir Bhalerao, Ligang He, Stephen A. Jarvis
Comput. J.8
2013 VSA: An offline scheduling analyzer for Xen virtual machine monitor
Zhiyuan Shao, Ligang He, Zhiqiang Lu, Hai Jin 0001
Future Gener. Comput. Syst.2
2013 Developing an optimized application hosting framework in Clouds
Xuanhua Shi, Hongbo Jiang 0001, Ligang He, Hai Jin 0001, Chonggang Wang, Xueguang Chen
J. Comput. Syst. Sci.3
2013 A DAG scheduling scheme on heterogeneous computing systems using double molecular structure-based chemical reaction optimization
Yuming Xu, Kenli Li 0001, Ligang He, Tung Khac Truong
J. Parallel Distributed Comput.3
2013 A Fast RPC System for Virtual Machines
abstract
Despite the advances in high performance interdomain communications for virtual machines (VM), data intensive applications developed for VMs based on the traditional remote procedure call (RPC) mechanism still suffer from performance degradation due to the inherent inefficiency of data serialization/deserilization operations. This paper presents VMRPC, a lightweight RPC framework specifically designed for VMs that leverages the heap and stack sharing mechanism to circumvent unnecessary data copy and serialization/deserilization. Our evaluation shows that the performance of VMRPC is an order of magnitude better than traditional RPC systems and existing alternative interdomain communication optimization systems. The evaluation on a VMRPC-enhanced networked file system across a varied range of benchmarks further reveals the competitiveness of VMRPC in IO-intensive applications.
Hao Chen 0002, Jianhua Sun 0002, Kenli Li 0001, Ligang He
IEEE Trans. Parallel Distributed Syst.5
2012 Performance Analysis for Workflow Management Systems under Role-Based Authorization Control
Ligang He, Stephen A. Jarvis
GPC2
2012 Chemical Reaction Optimization for Heterogeneous Computing Environments
abstract
Task scheduling has been proven to be NP-hard problem and we can usually approximate the best solutions with some classical algorithm, such as Heterogeneous Earliest Finish Time (HEFT), Genetic Algorithm. However, the huge types of scheduling problems and the small number of generally acknowledged methods mean that more methods are needed. In this paper, we propose a new method to schedule the execution of a group of dependent tasks for heterogeneous computing environments. The algorithm consists of two elements: An intelligent approach to assign the execution orders of tasks by task level, and an allocation algorithm based on chemical-reaction-inspired metaheuristic called Chemical Reaction Optimization (CRO) to map processors to tasks. The experiments show that the CRO-based algorithm performs consistently better than HEFT and Critical Path On a Processor (CPOP) without incurring much computational cost. Multiple runs of the algorithm can further improve the search result.
Kenli Li 0001, Yuming Xu, Bo Gao 0001, Ligang He
ISPA5
2012 Modeling and analyzing the impact of authorization on workflow executions
Ligang He, Chenlin Huang, Kewei Duan, Kenli Li 0001, Hao Chen 0002, Jianhua Sun 0002, Stephen A. Jarvis
Future Gener. Comput. Syst.1
2011 Dynamic Resource Allocation and Active Predictive Models for Enterprise Applications
Mohammad A. Alghamdi, Adam P. Chester, Ligang He, Stephen A. Jarvis
CLOSER3
2011 Modelling and analyzing the authorization and execution of video workflows
abstract
It is becoming common practice to migrate signal-based video workflows to IT-based Video workflows. Video workflows have some inherent features, including: 1) necessary human involvements in video workflows introduce security and authorization concerns; 2) the frequent change of video workflow contexts requires a flexible approach to acquiring performance data; 3) the content-centric nature of video workflows, which is in contrast to the business-centric of business workflows, requires the support of scheduled activities. This paper takes the above issues into account, proposing a novel mechanism for modeling video workflow executions in cluster-based resource pools under Role-Based Authorization Control (RBAC) schemes. The Color Timed Petri-Net (CTPN) formalism is applied to construct the models. Various types of authorization constraint are modeled in this paper, and scheduled activities are also supported in the model. There is a clear interface between workflow execution and workflow authorization modules. The constructed models are then simulated and analyzed to obtain performance data, including authorization overhead, system- and application-oriented performance. Based on the model analysis, this paper further proposes the methods to improve performance in the presence of authorization policies. This work can be used to plan system capacity subject to the authorization control, and can also be used to tune performance by changing the scheduling strategy and resource capacity when it is not possible to adjust the authorization policies.
Ligang He, Chenlin Huang, Kenli Li 0001, Hao Chen 0002, Jianhua Sun 0002, Bo Gao 0001, Kewei Duan, Stephen A. Jarvis
HiPC1
2011 Analyzing and Improving MPI Communication Performance in Overcommitted Virtualized Systems
abstract
Nowadays, it is an important trend in the system domain to use the software-based virtualization technology to build the execution environments (e.g., Clouds) and serve high performance computing (HPC) applications. However, with the extra virtualization layer, the application performance may be negatively affected. Studies revealed that the communication performance of the MPI library, which is widely used by the HPC applications, would suffer a high penalty when a physical host machine becomes overcommitted by virtual processors (VCPU). Unfortunately, the problem has not received enough attention and has not been solved yet in literature. In this paper, we investigate the reasons behind the performance penalty, and propose a solution to improve the communication performance of running MPI applications in the overcommitted virtualized systems. The experimental results show that by our proposal, most HPC applications can gain performance improvement to different extents among the overcommitted systems, depending on their communication patterns and the over committing level.
Zhiyuan Shao, Xuejiao Xie, Hai Jin 0001, Ligang He
MASCOTS5
2009 Performance prediction for running workflows under role-based authorization mechanisms
abstract
When investigating the performance of running scientific/commercial workflows in parallel and distributed systems, we often take into account only the resources allocated to the tasks constituting the workflow, assuming that computational resources will accept the tasks and execute them to completion once the processors are available. In reality, and in particular in Grid or e-business environments, security policies may be implemented in the individual organisations in which the computational resources reside. It is therefore expedient to have methods to calculate the performance of executing workflows under security policies. Authorisation control, which specifies who is allowed to perform which tasks when, is one of the most fundamental security considerations in distributed systems such as Grids. Role-Based Access Control (RBAC), under which the users are assigned to certain roles while the roles are associated with prescribed permissions, remains one of the most popular authorisation control mechanisms. This paper presents a mechanism to theoretically compute the performance of running scientific workflows under RBAC authorisation control. Various performance metrics are calculated, including both system-oriented metrics, (such as system utilisation, throughput and mean response time) and user-oriented metrics (such as mean response time of the workflows submitted by a particular client). With this work, if a client informs an organisation of the workflows they are going to submit, the organisation is able to predict the performance of these workflows running in its local computational resources (e.g. a high-performance cluster) enforced with RBAC authorisation control, and can also report client-oriented performance to each individual user.
Ligang He, Mark Calleja, Mark Hayes, Stephen A. Jarvis
IPDPS1
2009 CRBAC: Imposing multi-grained constraints on the RBAC model in the multi-application environment
Deqing Zou, Ligang He, Hai Jin 0001, Xueguang Chen
J. Netw. Comput. Appl.2
2008 Dynamic Resource Allocation in Enterprise Systems
abstract
It is common that Internet service hosting centres use several logical pools to assign server resources to different applications, and that they try to achieve the highest total revenue by making efficient use of these resources. In this paper, multi-tiered enterprise systems are modelled as multi-class closed queueing networks, with each network station corresponding to each application tier. In such queueing networks, bottlenecks can limit overall system performance, and thus should be avoided. We propose a bottleneck-aware server switching policy, which responds to system bottlenecks and switches servers to alleviate these problems as necessary. The switching engine compares the benefits and penalties of a potential switch, and makes a decision as to whether it is likely to be worthwhile switching. We also propose a simple admission control scheme, in addition to the switching policy, to deal with system overloading and optimise the total revenue of multiple applications in the hosting centre. Performance evaluation has been done via simulation and results are compared with those from a proportional switching policy and also a system that implements no switching policy. The experimental results show that the combination of the bottleneck-aware switching policy and the admission control scheme consistently outperforms the other two policies in terms of revenue contribution.
James Wen Jun Xue, Adam P. Chester, Ligang He, Stephen A. Jarvis
ICPADS3
2008 A System for Dynamic Server Allocation in Application Server Clusters
abstract
Application server clusters are often used to service high-throughput web applications. In order to host more than a single application, an organisation will usually procure a separate cluster for each application. Over time the utilisation of the clusters will vary, leading to variation in the response times experienced by users of the applications. Techniques that statically assign servers to each application prevent the system from adapting to changes in the workload, and are thus susceptible to providing unacceptable levels of service. This paper investigates a system for allocating server resources to applications dynamically, thus allowing applications to automatically adapt to variable workloads. Such a scheme requires meticulous system monitoring, a method for switching application servers between \text it {server pools} and a means of calculating when a server switch should be made (balancing switching cost against perceived benefits). Experimentation is performed using such a switching system on a Web application testbed hosting two applications across eight application servers. The testbed is used to compare several theoretically derived switching policies under a variety of workloads. Recommendations are made as to the suitability of different policies under different workload conditions.
Adam P. Chester, James Wen Jun Xue, Ligang He, Stephen A. Jarvis
ISPA3
2007 A scheduling algorithm for revenue maximisation for cluster-based Internet services
abstract
This paper proposes a new priority scheduling algorithm to maximise site revenue of session-based multi-tier Internet services in a multicluster environment. This research is part of a larger study in support of large-scale online trading systems and, as a result, this case study is chosen as a demonstrator for the techniques presented in this paper. The trading system is partitioned into a number of operations (trade, query etc.), which by their very nature are divided into orders of importance in terms of transactional response. The algorithm in this paper is based on Mean Value Analysis (MVA), which is used for the calculation of performance metrics concerning the queuing networks and workload allocation decision support in the multicluster. In addition to this, the priority assignment is based on combination of three attributes of any given request: (i) the sender class; (ii) the operation and, (iii) the status of the user’s portfolio (i.e.number of items in the user’s portfolio). A discrete event simulator has been developed to evaluate the performance of the priority scheduling scheme with different combinations of request attributes in various experimental scenarios. Our study aims to develop a dynamic scheduling policy, which takes into account real-time system parameters and optimises the site revenue. Although our priority scheduling algorithm is designed for an online trading system, it can be applied to most e-Commerce systems in which differentiated services are required.
James Wen Jun Xue, Ligang He, Stephen A. Jarvis
ICPADS2
2006 Performance evaluation of scheduling applications with DAG topologies on multiclusters with independent local schedulers
abstract
Before an application modelled as a directed acyclic graph (DAG) is executed on a heterogeneous system, a DAG mapping policy is often enacted. After mapping, the tasks (in the DAG-based application) to be executed at each computational resource are determined. The tasks are then sent to the corresponding resources, where they are orchestrated in the pre-designed pattern to complete the work. Most DAG mapping policies in the literature assume that each computational resource is a processing node of a single processor, i.e. the tasks mapped to a resource are to be run in sequence. Our studies demonstrate that if the resource is actually a cluster with multiple processing nodes, this assumption will cause a misperception in the tasks' execution time and execution order. This will disturb the pre-designed cooperation among tasks so that the expected performance cannot be achieved. In this paper, a DAG mapping algorithm is presented for multicluster architectures. Each constituent cluster in the multicluster is shared by background workload (from other users) and has its own independent local scheduler. The multicluster DAG mapping policy is based on theoretical analysis and its performance is evaluated through extensive experimental studies. The results show that compared with conventional DAG mapping policies, the new scheme that we present can significantly improve the scheduling performance of a DAG-based application in terms of the schedule length.
Ligang He, Stephen A. Jarvis, Daniel P. Spooner, Graham R. Nudd
IPDPS1
2006 Allocating Non-Real-Time and Soft Real-Time Jobs in Multiclusters
abstract
This paper addresses workload allocation techniques for two types of sequential jobs that might be found in multicluster systems, namely, non-real-time jobs and soft real-time jobs. Two workload allocation strategies, the optimized mean response time (ORT) and the optimized mean miss rate (OMR), are developed by establishing and numerically solving two optimization equation sets. The ORT strategy achieves an optimized mean response time for non-real-time jobs, while the OMR strategy obtains an optimized mean miss rate for soft real-time jobs over multiple clusters. Both strategies take into account average system behaviors (such as the mean arrival rate of jobs) in calculating the workload proportions for individual clusters and the workload allocation is updated dynamically when the change in the mean arrival rate reaches a certain threshold. The effectiveness of both strategies is demonstrated through theoretical analysis. These strategies are also evaluated through extensive experimental studies and the results show that when compared with traditional strategies, the proposed workload allocation schemes significantly improve the performance of job scheduling in multiclusters, both in terms of the mean response time (for non-real-time jobs) and the mean miss rate (for soft real-time jobs).
Ligang He, Stephen A. Jarvis, Daniel P. Spooner, Donna Dillenberger, Graham R. Nudd
IEEE Trans. Parallel Distributed Syst.1
2005 Mapping DAG-based applications to multiclusters with background workload
abstract
Before an application modelled as a directed acyclic graph (DAG) is executed on a heterogeneous system, a DAG mapping policy is often enacted. After mapping, the tasks (in the DAG-based application) to be executed at each computational resource are determined. The tasks are then sent to the corresponding resources, where they are orchestrated in the pre-designed pattern to complete the work. Most DAG mapping policies in the literature assume that each computational resource is a processing node of a single processor, i.e. the tasks mapped to a resource are to be run in sequence. Our studies demonstrate that if the resource is actually a cluster with multiple processing nodes, this assumption will cause a mis-perception in the tasks' execution time and execution order. This will disturb the pre-designed cooperation among tasks so that the expected performance cannot be achieved. In this paper, a DAG mapping algorithm is presented for multicluster architectures. Each constituent cluster in the multicluster is shared by background workload (from other users) and has its own independent local scheduler. The multicluster DAG mapping policy is based on theoretical analysis and its performance is evaluated through extensive experimental studies. The results show that compared with conventional DAG mapping policies, the new scheme that we present can significantly improve the scheduling performance of a DAG-based application in terms of the schedule length.
Ligang He, Stephen A. Jarvis, Daniel P. Spooner, David A. Bacigalupo, Guang Tan, Graham R. Nudd
CCGRID1
2005 Performance-Aware Workflow Management for Grid Computing
abstract
Grid middleware development has advanced rapidly over the past few years to support component-based programming models and service-oriented architectures. This is most evident with the forthcoming release of the Globus toolkit (GT4), which represents a convergence of concepts (and standards) from both the grid and web-services communities. Grid applications are increasingly modular, composed of workflow descriptions that feature both resource and application dynamism. Understanding the performance implications of scheduling grid workflows is critical in providing effective resource management and reliable service quality to users. This paper describes a series of extensions to an existing performance-aware grid management system (TITAN). These extensions provide additional support for workflow prediction and scheduling using a multi-domain performance management infrastructure.
Daniel P. Spooner, Stephen A. Jarvis, Ligang He, Graham R. Nudd
Comput. J.4
2005 The impact of predictive inaccuracies on execution scheduling
Stephen A. Jarvis, Ligang He, Daniel P. Spooner, Graham R. Nudd
Perform. Evaluation2
2005 An Investigation into the Application of Different Performance Prediction Methods to Distributed Enterprise Applications
David A. Bacigalupo, Stephen A. Jarvis, Ligang He, Daniel P. Spooner, Donna Dillenberger, Graham R. Nudd
J. Supercomput.3
2004 An Investigation into the Application of Different Performance Prediction Techniques to e-Commerce Applications
abstract
Summary form only given. Predictive performance models of e-Commerce applications allows grid workload managers to provide e-Commerce clients with qualities of service (QoS) whilst making efficient use of resources. We demonstrate the use of two 'coarse-grained' modelling approaches (based on layered queuing modelling and historical performance data analysis) for predicting the performance of dynamic e-Commerce systems on heterogeneous servers. Results for a popular e-Commerce benchmark show how request response times and server throughputs can be predicted on servers with heterogeneous CPUs at different background loads. The two approaches are compared and their usefulness to grid workload management is considered.
David A. Bacigalupo, Stephen A. Jarvis, Ligang He, Graham R. Nudd
IPDPS3
2004 Optimising Static Workload Allocation in Multiclusters
abstract
Summary form only given. Workload allocation and job dispatching are two fundamental components in static job scheduling for distributed systems. We address the static workload allocation techniques for two types of job stream in multicluster systems, namely, nonreal-time job streams and soft-real-time job streams, which request different qualities of service. Two workload allocation strategies (called ORT and OMR) are developed by establishing and numerically solving two optimisation equation sets. The ORT strategy achieves the optimised mean response time for the nonreal-time job stream; while the OMR strategy can gain the optimised mean miss rate for the soft-real-time job stream over multiple clusters (these strategies can also be applied in a single cluster system). The effectiveness of both strategies is demonstrated through theoretical analysis. The proposed workload allocation schemes are combined with two job dispatching strategies (weighted random and weighted round-robin) to generate new static job scheduling algorithms for multicluster environments. These algorithms are evaluated through extensive experimental studies and the results show that compared with static approaches without the optimisation techniques, the proposed workload allocation schemes can significantly improve the performance of static job scheduling in multiclusters, in terms of both the mean response time (for the nonreal-time jobs) and the mean miss rate (for soft-real-time jobs).
Ligang He, Stephen A. Jarvis, Daniel P. Spooner, Graham R. Nudd
IPDPS1
2004 Performance-Aware Load Balancing for Multiclusters
Ligang He, Stephen A. Jarvis, David A. Bacigalupo, Daniel P. Spooner, Graham R. Nudd
ISPA1
2003 Dynamic Scheduling of Parallel Real-Time Jobs by Modelling Spare Capabilities in Heterogeneous Clusters
abstract
In this research, a scenario is assumed where periodic real-time jobs are being run on a heterogeneous cluster of computers, and new aperiodic parallel real-time jobs, modelled by directed acyclic graphs (DAG), arrive at the system dynamically. In the scheduling scheme presented in this paper, a global scheduler situated within the cluster schedules new jobs onto the computers by modelling their spare capabilities left by existing periodic jobs. Admission control is introduced so that new jobs are rejected if their deadlines cannot be met under the precondition of still guaranteeing the real-time requirements of existing jobs. Each computer within the cluster houses a local scheduler, which uniformly schedules both periodic job instances and the subtasks in the parallel realtime jobs using an early deadline first policy. The modelling of the spare capabilities is optimal in the sense that once a new task starts running on a computer, it will utilize all the spare capability left by the periodic real-time jobs and its finish time is the earliest possible. The performance of the proposed modelling approach and scheduling scheme is evaluated by extensive simulation; results show that the system utilization is significantly enhanced, while the real-time requirements of the existing jobs remain guaranteed.
Ligang He, Stephen A. Jarvis, Daniel P. Spooner, Graham R. Nudd
CLUSTER1
2003 Performance-Based Dynamic Scheduling of Hybrid Real-Time Applications on a Cluster of Hetrogeneous Workstations
Ligang He, Stephen A. Jarvis, Daniel P. Spooner, Graham R. Nudd
Euro-Par1
2001 Optimal Scheduling of Aperiodic Jobs on Cluster
Ligang He, Hai Jin 0001, Zongfen Han
Euro-Par1