Cho-Li Wang

dblp:44/3271 · DBLP profile ↗
← Back
119ranked-venue papers
3as first author
15since 2021 · last 2025
0000-0002-4629-7175ORCID · corroborated

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

Systems, architecture and hardware · 89 · 2 first-author · 15 since 2021Computer networks · 6Human-computer interaction and ubiquitous computing · 4Applied, interdisciplinary, general and emerging computing · 4 · 1 first-authorArtificial intelligence and machine learning · 3Software engineering, systems software and programming languages · 2Security and privacy · 1Databases, data management, data science and information retrieval · 1Graphics, computer vision, multimedia, augmented reality and games · 1
YearPublicationVenuePosition
2025 Cube-fx: Mapping Taylor Expansion Onto Matrix Multiplier-Accumulators of Huawei Ascend AI Processors
abstract
Taylor expansion, a mature method for function evaluations used in Artificial Intelligence (AI) applications, approximates functions with polynomials. In addition to the function evaluations, AI applications require massive matrix multiplications, inspiring manufacturers to propose AI processors with matrix multiplier-accumulators (MACs). However, compared with the powerful Matrix MACs, the vectorized units of the AI processors cannot efficiently carry the existing Taylor expansion implementation of Single Instruction Multiple Data (SIMD) parallelism. Leveraging the Matrix MACs for Taylor expansion becomes an ideal direction. In previous studies, migrating optimized algorithms to the Matrix MACs requires matrix generation during the runtime. The generation is expensive and even cancels the accelerations brought by the Matrix MACs on the AI processors, which Taylor expansion also suffers. This article presents Cube-fx, a mapping algorithm of Taylor expansion for multiple functions onto Matrix MACs. Cube-fx expresses the building and computation in matrix multiplications without inefficient dynamic matrix generation. On Huawei Ascend processors, Cube-fx averagely achieves 1.64× speedups compared with vectorized Horner's Method with 56.38$\%$vectorized operations reduced.
Yifeng Tang, Huaman Zhou, Zhuoran Ji, Cho-Li Wang
IEEE Trans. Parallel Distributed Syst.4
2023 Embedding Communication for Federated Graph Neural Networks with Privacy Guarantees
abstract
Graph Neural Networks (GNNs) have been widely used in many Machine Learning (ML) tasks as they show remarkable performance in modeling graph structure data. While several distributed GNN frameworks have been proposed to tackle the training of huge graphs, uploading the local graph data to the central server for model training is impractical in real-world scenarios due to privacy concerns. Federated Learning (FL) is introduced as an effective technology to address the privacy issue, allowing edge clients to collaboratively train the ML models locally. However, GNNs follow a recursive neighborhood aggregation scheme. Computing the representation vector, also known as embedding, of one node requires aggregating feature vectors of its neighbors. Training the model based on local subgraphs would suffer from information loss and result in accuracy degradation. This paper presents EmbC-FGNN, an efficient Federated _Graph Neural Network framework that enables Node Embedding Communication among training clients in a privacy-preserving way. EmbC-FGNN first proposes an Embedding Server (ES) to maintain and synchronize the shared embeddings among edge workers. It allows training devices to expand the local subgraphs with exchanged embeddings to improve the model accuracy without revealing local node features and graph topology. To minimize the communication costs of the ES, we introduce a periodic embedding synchronization strategy to reduce the communication frequency. Furthermore, we apply asynchronous training to accelerate the convergence speed. Experimental results on several graph neural networks and datasets demonstrate that EmbC-FGNN can improve the overall accuracy (more than 10% for Reddit dataset) and achieve good round-to-accuracy performance.
Xueyu Wu 0001, Zhuoran Ji, Cho-Li Wang
ICDCS3
2023 SelB-k-NN: A Mini-Batch K-Nearest Neighbors Algorithm on AI Processors
abstract
The popularity of Artificial Intelligence (AI) motivates novel domain-specific hardware named AI processors. With a design trade-off, the AI processors feature incredible computation power for matrix multiplications and activations, while some leave other operations less powerful, e.g., scalar operations and vectorized comparisons & selections. For k-nearest neighbors (k-NN) algorithm, consisting of distance computation phase and k-selection phase, while the former is naturally accelerated, the previous efficient k-selection becomes problematic. Moreover, limited memory forces k-NN to adopt a mini-batch manner with tiling technique. As the distance computation’s results are the k-selection’s inputs, the former’s tiling shape determines that of the latter. Since the two phases execute on separate hardware units requiring different performance analyses, whether the former’s tiling strategies benefit the latter and entire k-NN is doubtful.To address the new challenges brought by the AI processors, this paper proposes SelB-k-NN (Selection-Bitonic-k-NN), a mini-batch algorithm inspired by selection sort and bitonic k-selection. SelB-k-NN avoids the expansion of the weakly-supported operations on the huge scale of datasets. To apply SelB-k-NN to various AI processors, we propose two algorithms to reduce the hardware support requirements. Since the matrix multiplication operates data with the specifically-designed memory hierarchy which k-selection does not share, the tiling shape of the former cannot guarantee the best execution of the latter and vice versa. By quantifying the runtime workload variations of k-selection, we formulate an optimization problem to search for the optimal tiling shapes of both phases with an offline pruning method, which reduces the search space in the preprocessing stage. Evaluations show that on Huawei Ascend 310 AI processor, SelB-k-NN achieves 2.01× speedup of the bitonic k-selection, 23.93× of the heap approach, 78.52× of the CPU approach. For mini-batch SelB-k-NN, the optimal tiling shapes for two phases respectively achieve 1.48× acceleration compared with the matrix multiplication tiling shapes and 1.14× with the k-selection tiling shapes, with 72.80% of the search space pruned.
Yifeng Tang, Cho-Li Wang
IPDPS2
2023 Performance modeling on DaVinci AI core
Yifeng Tang, Cho-Li Wang
J. Parallel Distributed Comput.2
2022 Optimizing Aggregate Computation of Graph Neural Networks with on-GPU Interpreter-Style Programming
abstract
Graph Neural Networks (GNNs) generalize deep learning to graph-structured data and show great success in many tasks. However, their irregular aggregation kernels make them inefficient on GPUs. The unpredictable control flow and memory references of irregular kernels prohibit most optimizations designed for regular ones. For example, even if the nodes have overlapped neighbors, reusing them via shared memory is non-trivial, as the neighborhoods used are runtime information. This paper presents regGNN, an aggregation implementation that can benefit from the optimizations designed for regular kernels. It proposes a concept named "semi-regular" to describe the aggregate computation: the irregularity only comes from the neighborhood traversal; aggregating the high-dimensional vectors, which dominates the computation, is data-independent and thus incurs no irregularity. regGNN encodes the aggregate computation steps of each thread block into an aggregate script, which replaces the graph as an input of the GPU kernel. The GPU kernel is like an interpreter, and the aggregate script can be regarded as written in a simple GPU scripting language. The optimizations designed for regular kernels can then be applied to the aggregate script, as it is static and regular. regGNN demonstrates three optimizations: (1) intelligently scheduling nodes and customizing shared memory replacement to maximize data reuse, (2) reassigning nodes among warps for load balancing, and (3) aligning the aggregate script to improve memory latency hiding. Compared with the state-of-the-art GNN frameworks, regGNN achieves 2.81× throughput on average for moderate-scale GNNs. The speedup increases to 5.21× for GNNs with small hidden sizes and 100s × for deep GNNs.
Zhuoran Ji, Cho-Li Wang
PACT2
2022 KAFL: Achieving High Training Efficiency for Fast-K Asynchronous Federated Learning
abstract
Federated Averaging (FedAvg) and its variants are prevalent optimization algorithms adopted in Federated Learning (FL) as they show good model convergence. However, such optimization methods are mostly running in a synchronous flavor which is plagued by the straggler problem, especially in the real-world FL scenario. Federated learning involves a massive number of resource-weak edge devices connected to the intermittent networks, exhibiting a vastly heterogeneous training environment. The asynchronous setting is a plausible solution to fulfill the resources utilization. Yet, due to data and device heterogeneity, the training bias and model staleness dramatically downgrade the model performance. This paper presents KAFL, a fast-K Asynchronous Federated Learning framework, to improve the system and statistical efficiency. KAFL allows the global server to iteratively collect and aggregate (1) the parameters uploaded by the fastest K edge clients (K-FedAsync); or (2) the first M updated parameters sent from any clients (Mstep-FedAsync). Compared to the fully asynchronous setting, KAFL helps the server obtain a better direction toward the global optima as it collects the information from at least K clients or M parameters. To further improve the convergence speed of KAFL, we propose a new weighted aggregation method which dynamically adjusts the aggregation weights according to the weight deviation matrix and client contribution frequency. Experimental results show that KAFL achieves a significant time-to-target-accuracy speedup on both IID and Non-IID datasets. To achieve the same model accuracy, KAFL reduces more than 50% training time for five CNN and RNN models, demonstrating the high training efficiency of our proposed framework.
Xueyu Wu 0001, Cho-Li Wang
ICDCS2
2022 Efficient exact K-nearest neighbor graph construction for billion-scale datasets using GPUs with tensor cores
abstract
Approximate nearest neighbor search plays a fundamental role in many areas, and the k-nearest neighbor graph (KNNG) becomes a promising solution, especially in high-dimensional space. The advantages of KNNG come at the expense of high construction time, which is in quadratic time complexity in the number of points. Many GPUs have adopted specialized hardware units for matrix multiplication, providing an even higher arithmetic throughput. This paper presents flyKNNG, a GPU KNNG construction algorithm for billion-scale datasets. It deploys the distance matrix calculation to matrix multiplication units and adopts on-the-fly top-k selection to avoid transferring the exa-scale distance matrix to/from device memory. flyKNNG co-designs the two key algorithms to optimize the overall performance: the distance matrix calculation algorithm considers the data communication costs and pruning strategy of top-k selection; the top-k selection algorithm is also specially designed for on-the-fly usage, which impacts the data reuse and instruction-level parallelism of the distance matrix calculation as little as possible. Moreover, our top-k selection algorithm is optimized for the special data layout adopted by most matrix multiplication units. Experiments show that flyKNNG achieves 4.67X speedup compared with CUML/FAISS, one of the state-of-the-art approaches.
Zhuoran Ji, Cho-Li Wang
ICS2
2022 Compiler-Directed Incremental Checkpointing for Low Latency GPU Preemption
abstract
GPUs are widely used in data centers to accelerate data-parallel applications. The multiuser and multitasking environment provides a strong incentive for preemptive GPU multitasking, especially for latency-sensitive jobs. Due to the large contexts of GPU kernels, preemptive GPU context switching is costly. Many novel GPU preemption techniques are proposed. Among them, checkpoint-based GPU preemption enables low latency GPU preemption but incurs a high runtime overhead. Prior studies propose to exclude dead registers from the checkpoint file to reduce the runtime overhead. It works well for CPUs, but it is not rare that a live register is not updated between two checkpoints for GPU kernels. This paper presents TripleC, a compiler-directed incremental checkpointing technique specially designed for GPU preemption. It further excludes the registers, which have not been overwritten since the last time they were spilled, from the checkpoint file with data flow analysis. The checkpoint placement algorithm of TripleC can properly estimate a checkpoint's cost under incremental checkpointing. It also considers the interaction among checkpoints so that the overall cost is minimized. Moreover, TripleC relaxes the conventional checkpointing constraint that the whole register context must be spilled before passing the checkpoint. Because of the diverse control flow, placing a register spilling instruction at different points incurs different costs. TripleC minimizes the cost with a two-phase algorithm that schedules these register spilling instructions at compilation time. Evaluations show that TripleC reduces the runtime overhead by 12.9 % on average compared with the state-of-the-art non-incremental checkpointing approach.
Zhuoran Ji, Cho-Li Wang
IPDPS2
2022 Momentum-driven adaptive synchronization model for distributed DNN training on HPC clusters
Zhuoran Ji, Cho-Li Wang
J. Parallel Distributed Comput.3
2022 SaPus: Self-Adaptive Parameter Update Strategy for DNN Training on Multi-GPU Clusters
abstract
Parameter server architecture has been identified as an efficient framework for scaling DNNs training on clusters. For large-scale deployment, communication becomes the bottleneck, and the parameter updating strategy strongly impacts the training performance and accuracy. Recent state-of-art solutions have adopted the local SGD approach, which enables workers to update their local version of models and only aggregate them to update the global parameters after finishing a number of iterations, to alleviate heavy communication pressure on the parameter server and improving the training performance. We identify three limitations of these works. First, these works do not provide an approach for determining when the worker is to update the parameter with the server under asynchronous communication strategies that can guarantee the training performance. Second, local SGD suffers from the problem of unbounded gradient delay. Previous work works well for a short delay while can not guarantee the performance with an increase of gradient delay. Third, they do not consider the system performance when determining the update interval of the local SGD, including the CPU, memory, and network, which affects the training performance extremely. We provide a self-adaptive parameter updating strategy called SaPus, which allows each worker to detect their training results through quantification of the accumulated gradient updates and determine when to update the parameter with the server adaptively and individually. Theoretical lower and upper bound of the update interval is also provided. We also propose a weighted aggregation algorithm based on a global-loss window, which is used to collect the most recent loss value of other workers to calculate aweightfor the accumulated gradients of each worker to solve the unbounded delay problem in asynchronous local SGD. To increase the robustness of our parameter updating strategy, a performance model is built to provide a resource-aware lower bound for the update interval. Extensive experimental results generated on GPU cluster indicate that our model improves the training performance of DNNs, achieving up to$66.67\%$speedup as compared with state-of-art solutions. Further, results show the CPU utilization of server dropped by up to$81.1\%$and network bandwidth usage reduced to less than$1~Gbps$on an average during the training.
Cho-Li Wang
IEEE Trans. Parallel Distributed Syst.2
2022 MIPD: An Adaptive Gradient Sparsification Framework for Distributed DNNs Training
abstract
Asynchronous training based on the parameter server architecture is widely used for scaling up the DNN training over large datasets and DNN models. Communication has been identified as the major bottleneck when deploying the DNN training over the large-scale distributed deep learning systems. Recent studies try to reduce the communication traffic through gradient sparsification and quantization approaches. We identify three limitations in previous studies. First, the fundamental guideline for gradient sparsification of their work is the magnitude of the gradient. However, the gradients’ magnitude represents the current optimization direction while it cannot indicate the significance of the parameters, which potentially results in delayed updating for the significant parameters. Second, their gradient quantization methods based on the entire model often lead to error accumulation for gradients aggregation since the gradients from different layers of the DNN model follow different distributions. Third, previous quantization approaches are CPU intensive, which generates strong overhead for the server. We proposeMIPD, an adaptive and layer-wise gradient sparsification framework that compresses the gradients based on model interpretability and probability distribution of gradients. MIPD compresses the gradients according to the corresponding significance of its parameters, which is defined by model interpretability. An Exponential Smoothing method is also proposed to compensate for the dropped gradients on the server to reduce the gradients error. MIPD proposes to update half of the parameters for each training step to reduce the CPU overhead of the server. It encodes the gradients based on their probability distribution, thereby minimizing the approximated errors. Extensive experimental results generated on the GPU cluster indicate that the proposed framework effectively improves the training performance of DNNs by up to 36.2%, which ensures high accuracy as compared to state-of-art solutions. Accordingly, the CPU and network usage of the server dropped by up to 42.0% and 32.7% respectively.
Cho-Li Wang
IEEE Trans. Parallel Distributed Syst.2
2021 Collaborative GPU Preemption via Spatial Multitasking for Efficient GPU Sharing
Zhuoran Ji, Cho-Li Wang
Euro-Par2
2021 Accelerating DBSCAN Algorithm with AI Chips for Large Datasets
abstract
DBSCAN is a popular clustering algorithm, which shows great success in many real-world applications. Its advantages come at the expense of massive computation, especially for computing the distance matrix. Driven by deep learning, many Artificial Intelligence (AI) chips have been developed. With efficient matrix multiplication units, AI chips can significantly accelerate the distance calculation. However, DBSCAN also needs to identify and count the neighbors for each point. It is challenging for most AI chips due to over-specialization. Moreover, the increasing data size and the limited device memory capacity force DBSCAN to follow a mini-batch manner. It results in a high data transfer overhead, which further hinders the performance of DBSCAN on AI chips. In this paper, we propose two novel techniques to address the challenges of accelerating the DBSCAN algorithm with AI chips: (1) new neighbor identification algorithms using bitwise operations only, while traditional solutions require the compare-and-select operations that are weakly supported in AI chips; and (2) two speculative execution strategies to reduce the data transfer overhead induced by mini-batches. Evaluations show that deploying distance matrix calculation to tensor cores achieves 2.61 × speedup on Nvidia RTX 3090. On Huawei Ascend 310, our neighbor identification algorithms achieve 17.88 × throughout of using CPUs for neighbor identification. The speculative execution strategies further reduce the execution time by 15.1% on average for normal datasets and up to 99.0% for sparse datasets.
Zhuoran Ji, Cho-Li Wang
ICPP2
2021 CTXBack: Enabling Low Latency GPU Context Switching via Context Flashback
abstract
Efficient GPU preemption mechanisms are critical for task prioritization in multitasking environments, especially for latency-sensitive GPU applications. However, due to the large context of GPU kernels, simply borrowing context switching mechanisms from CPU space incurs substantial latency and overhead. To address this problem, we propose CTXBack, which allows a thread block to execute context switching at a preceding instruction with a smaller context. It enables low latency GPU context switching for latency-sensitive applications on shared GPUs. Three complementary ways are proposed to make more preceding instructions into valid points to execute context switching, giving CTXBack a higher chance to find instructions with smaller contexts. Evaluations show that CTXBack reduces the context by 61.0%, which is only 1.09× of the minimum possible context size. With only 0.41% runtime overhead, the preemption latency and resuming time are reduced by 63.1% and 50.0% on average compared to the traditional approach.
Zhuoran Ji, Cho-Li Wang
IPDPS2
2021 FedSCR: Structure-Based Communication Reduction for Federated Learning
abstract
Federated Learning allows edge devices to collaboratively train a shared model on their local data without leaking user privacy. The non-independent-and-identically-distributed (Non-IID) property of data distribution, which leads to severe accuracy degradation, and enormous communication overhead for aggregating parameters should be tackled in federated learning. In this article, we conduct a detailed analysis of parameter updates on the Non-IID datasets and compare the difference with the IID setting. Experimental results exhibit that parameter update matrices are structure-sparse and show that more gradients could be identified as negligible updates on the Non-IID data. As a result, we propose a structure-based communication reduction algorithm, called FedSCR, that reduces the number of parameters transported through the network while maintaining the model accuracy. FedSCR aggregates the parameter updates over channels and filters, identifies and removes the redundant updates by comparing the aggregated values with a threshold. Unlike the traditional structured pruning methods, FedSCR retains the complete model that does not require to be retrained and fine-tuned. The local loss and weight divergence on each device vary a lot because of the unbalanced data distribution. We further propose an adaptive FedSCR, that dynamically changes the bounded threshold, to enhance the model robustness on the Non-IID data. Evaluation results show that our proposed strategies achieve almost 50 percent upstream communication reduction without loss of accuracy. FedSCR can be integrated into state-of-the-art federated learning algorithms to dramatically reduce the number of parameters pushed to the global server with a tolerable accuracy reduction.
Xueyu Wu 0001, Xin Yao 0008, Cho-Li Wang
IEEE Trans. Parallel Distributed Syst.3
2020 Uranus: Simple, Efficient SGX Programming and its Applications
abstract
Applications written in Java have strengths to tackle diverse threats in public clouds, but these applications are still prone to privileged attacks when processing plaintext data. Intel SGX is powerful to tackle these attacks, and traditional SGX systems rewrite a Java application's sensitive functions, which process plaintext data, using C/C++ SGX API. Although this code-rewrite approach achieves good efficiency and a small TCB, it requires SGX expert knowledge and can be tedious and error-prone. To tackle the limitations of rewriting Java to C/C++, recent SGX systems propose a code-reuse approach, which runs a default JVM in an SGX enclave to execute the sensitive Java functions. However, both recent study and this paper find that running a default JVM in enclaves incurs two major vulnerabilities, Iago attacks, and control flow leakage of sensitive functions, due to the usage of OS features in JVM. In this paper, Uranus creates easy-to-use Java programming abstractions for application developers to annotate sensitive functions, and Uranus automatically runs these functions in SGX at runtime. Uranus effectively tackles the two major vulnerabilities in the code-reuse approach by presenting two new protocols: 1) a Java bytecode attestation protocol for dynamically loaded functions; and 2) an OS-decoupled, efficient GC protocol optimized for data-handling applications running in enclaves. We implemented Uranus in Linux and applied it to two diverse data-handling applications: Spark and ZooKeeper. Evaluation shows that: 1) Uranus achieves the same security guarantees as two relevant SGX systems for these two applications with only a few annotations; 2) Uranus has reasonable performance overhead compared to the native, insecure applications; and 3) Uranus defends against privileged attacks. Uranus source code and evaluation results are released on https://github.com/hku-systems/uranus.
Jianyu Jiang, Xusheng Chen, Tsz On Li, Cheng Wang 0021, Tianxiang Shen, Shixiong Zhao, Heming Cui, Cho-Li Wang, Fengwei Zhang
AsiaCCS8
2020 On-GPU thread-data remapping for nested branch divergence
Huanxin Lin, Cho-Li Wang
J. Parallel Distributed Comput.2
2020 A Model-Based Software Solution for Simultaneous Multiple Kernels on GPUs
abstract
As a critical computing resource in multiuser systems such as supercomputers, data centers, and cloud services, a GPU contains multiple compute units (CUs). GPU Multitasking is an intuitive solution to underutilization in GPGPU computing. Recently proposed solutions of multitasking GPUs can be classified into two categories: (1) spatially partitioned sharing (SPS), which coexecutes different kernels on disjointed sets of compute units (CU), and (2) simultaneous multikernel (SMK), which runs multiple kernels simultaneously within a CU. Compared to SPS, SMK can improve resource utilization even further due to the interleaving of instructions from kernels with low dynamic resource contentions. However, it is hard to implement SMK on current GPU architecture, because (1) techniques for applying SMK on top of GPU hardware scheduling policy are scarce and (2) finding an efficient SMK scheme is difficult due to the complex interferences of concurrently executed kernels. In this article, we propose a lightweight and effective performance model to evaluate the complex interferences of SMK. Based on the probability of independent events, our performance model is built from a totally new angle and contains limited parameters. Then, we propose a metric, symbiotic factor , which can evaluate an SMK scheme so that kernels with complementary resource utilization can corun within a CU. Also, we analyze the advantages and disadvantages of kernel slicing and kernel stretching techniques and integrate them to apply SMK on GPUs instead of simulators. We validate our model on 18 benchmarks. Compared to the optimized hardware-based concurrent kernel execution whose kernel launching order brings fast execution time, the results of corunning kernel pairs show 11%, 18%, and 12% speedup on AMD R9 290X, RX 480, and Vega 64, respectively, on average. Compared to the Warped-Slicer, the results show 29%, 18%, and 51% speedup on AMD R9 290X, RX 480, and Vega 64, respectively, on average.
Hao Wu 0130, Huanxin Lin, Cho-Li Wang
ACM Trans. Archit. Code Optim.4
2020 Probabilistic Consistency Guarantee in Partial Quorum-Based Data Store
abstract
Many NoSQL databases support quorum-based protocols, which require a subset of replicas (called a quorum) to respond to each write/read operation. These systems configure the quorum size to tune the operation latency and adopt multiple consistency levels. Some recent works illustrate that using probability models to quantify the chance of reading the last update is important because it could avoid returning stale values under eventual consistency. There are two challenging issues: (1) from inconsistent replicas, how to determine the minimum quorum size (i.e., the lowest access latency) to read the newest data at a specified probability; (2) node failure frequently happens in large-scale systems, how to guarantee the probability-based consistent reads. This article presents Probabilistic Consistency Guarantee (PCG), which is the first dynamic quorum decision and failure-aware quantification model. PCG model respectively quantifies the server-side consistency after the latest write, which reflects the object's time-varying update progress, and the possibility of reading this update when responding to the end-users. Our theoretical analysis derives several formulas to determine the quorum size of a read quorum and the consensus result selected from this quorum is the data updated by the last write at the user-specified probability. When some replicas are unavailable, our model knows how to rescale the quorum and read values from surviving replicas could reduce the stale reads caused by node failures. The experimental results in Cassandra demonstrate that the PCG model can achieve up to 77.7 percent more accurate predictions and reduce up to 48.9 percent read latency than those of the previous model.
Xin Yao 0008, Cho-Li Wang
IEEE Trans. Parallel Distributed Syst.2
2019 EC-Shuffle: Dynamic Erasure Coding Optimization for Efficient and Reliable Shuffle in Spark
abstract
Fault-tolerance capabilities attract increasing attention from existing data processing frameworks, such as Apache Spark. To avoid replaying costly distributed computation, like shuffle, local checkpoint and remote replication are two popular approaches. They incur significant runtime overhead, such as extra storage cost or network traffic. Erasure coding is another emerging technology, which also enables data resilience. It is perceived as capable of replacing the checkpoint and replication mechanisms for its high storage efficiency. However, it suffers heavy network traffic due to distributing data partitions to different locations. In this paper, we propose EC-Shuffle with two encoding schemes and optimize the shuffle-based operations in Spark or MapReduce-like frameworks. Specifically, our encoding schemes concentrate on optimizing the data traffic during the execution of shuffle operations. They only transfer the parity chunks generated via erasure coding, instead of a whole copy of all data chunks. EC-Shuffle also provides a strategy, which can dynamically select the per-shuffle biased encoding scheme according to the number of senders and receivers in each shuffle. Our analyses indicate that this dynamic encoding selection can minimize the total size of parity chunks. The extensive experimental results using BigDataBench with hundreds of mappers and reducers shows this optimization can reduce up to 50% network traffic and achieve up to 38% performance improvement.
Xin Yao 0008, Cho-Li Wang, Mingzhe Zhang 0003
CCGRID2
2019 FluentPS: A Parameter Server Design with Low-frequency Synchronization for Distributed Deep Learning
abstract
With pursuing high accuracy on big datasets, current research prefers designing complex neural networks, which need to maximize data parallelism for short training time. Many distributed deep learning systems, such as MXNet and Petuum, widely use parameter server framework with relaxed synchronization models. Although these models could cost less on each synchronization, its frequency is still high among many workers, e.g., the soft barrier introduced by Stale Synchronous Parallel (SSP) model. In this paper, we introduce our parameter server design, namely FluentPS, which can reduce frequent synchronization and optimize communication overhead in a large-scale cluster. Different from using a single scheduler to manage all parameters' synchronization in some previous designs, our system allows each server to independently adjust schemes for synchronizing its parameter shard and overlaps the push and pull processes of different servers. We also explore two methods to improve the SSP model: (1) lazy execution of buffered pull requests to reduce the synchronization frequency and (2) a probability-based strategy to pause the fast worker at a probability under SSP condition, which avoids unnecessary waiting of fast workers. We evaluate ResNet-56 with the same large batch size at different cluster scales. While guaranteeing robust convergence, FluentPS gains up to 6× speedup and reduce 93.7% communication time costs than PS-Lite. The raw SSP model causes up to 131× delayed pull requests than our improved synchronization model, which can provide fine-tuned staleness controls and achieve higher accuracy.
Xin Yao 0008, Xueyu Wu 0001, Cho-Li Wang
CLUSTER3
2019 Spectral Graph Theory Based Topology Analysis for Reconfigurable Data Center Networks
abstract
Emerging technological innovations introduce the possibility to reconfigure the data center topology at runtime. The development of reconfigurable architectures can adapt their topology to account for changing demands, e.g., using Flyways-like augmented links. However, there is no common notion established in this area of how to evaluate the topology, the underlying theoretical analysis is not yet well studied. In this paper, we present the upper bound of the network diameter using spectral graph theory, to the best of our knowledge, which is a first theoretical attempt on understanding the nature of the reconfigurable data center networks. We further prove the algebraic connectivity is insensitive to the weight changes. Finally, based the algebra-connectivity λ2, we give a comprehensive link-augmentation validity, which can be implemented in current DCNs potentially.
Dengke Zhang, Xingwei Wang 0001, Min Huang 0001, Cho-Li Wang
MSN4
2019 Efficient low-latency packet processing using On-GPU Thread-Data Remapping
Huanxin Lin, Cho-Li Wang
J. Parallel Distributed Comput.2
2018 On-GPU Thread-Data Remapping for Branch Divergence Reduction
abstract
General Purpose GPU computing (GPGPU) plays an increasingly vital role in high performance computing and other areas like deep learning. However, arising from the SIMD execution model, the branch divergence issue lowers efficiency of conditional branching on GPUs, and hinders the development of GPGPU. To achieve runtime on-the-spot branch divergence reduction, we propose the first on-GPU thread-data remapping scheme. Before kernel launching, our solution inserts codes into GPU kernels immediately before each target branch so as to acquire actual runtime divergence information. GPU software threads can be remapped to datasets multiple times during single kernel execution. We propose two thread-data remapping algorithms that are tailored to the GPU architecture. Effective on two generations of GPUs from both NVIDIA and AMD, our solution achieves speedups up to 2.718 with third-party benchmarks. We also implement three GPGPU frontier benchmarks from areas including computer vision, algorithmic trading and data analytics. They are hindered by more complex divergence coupled with different memory access patterns, and our solution works better than the traditional thread-data remapping scheme in all cases. As a compiler-assisted runtime solution, it can better reduce divergence for divergent applications that gain little acceleration on GPUs for the time being.
Huanxin Lin, Cho-Li Wang, Hongyuan Liu 0002
ACM Trans. Archit. Code Optim.2
2018 SIMPO: A Scalable In-Memory Persistent Object Framework Using NVRAM for Reliable Big Data Computing
abstract
While CPU architectures are incorporating many more cores to meet ever-bigger workloads, advance in fault-tolerance support is indispensable for sustaining system performance under reliability constraints. Emerging non-volatile memory technologies are yielding fast, dense, and energy-efficient NVRAM that can dethrone SSD drives for persisting data. Research on using NVRAM to enable fast in-memory data persistence is ongoing. In this work, we design and implement a persistent object framework, dubbed scalable in-memory persistent object (SIMPO) , which exploits NVRAM, alongside DRAM, to support efficient object persistence in highly threaded big data applications. Based on operation logging, we propose a new programming model that classifies functions into instant and deferrable groups. SIMPO features a streamlined execution model, which allows lazy evaluation of deferrable functions and is well suited to big data computing workloads that would see improved data locality and concurrency. Our log recording and checkpointing scheme is effectively optimized towards NVRAM, mitigating its long write latency through write-combining and consolidated flushing techniques. Efficient persistent object management with features including safe references and memory leak prevention is also implemented and tailored to NVRAM. We evaluate a wide range of SIMPO-enabled applications with machine learning, high-performance computing, and database workloads on an emulated hybrid memory architecture and a real hybrid memory machine with NVDIMM. Compared with native applications without persistence, experimental results show that SIMPO incurs less than 5% runtime overhead on both platforms and even gains up to 2.5× speedup and 84% increase in throughput in highly threaded situations on the two platforms, respectively, thanks to the streamlined execution model.
Mingzhe Zhang 0003, King Tin Lam, Xin Yao 0008, Cho-Li Wang
ACM Trans. Archit. Code Optim.4
2018 Confluence: Speeding Up Iterative Distributed Operations by Key-Dependency-Aware Partitioning
abstract
A typical shuffle operation randomly partitions data on many computers, generating possibly a significant amount of network traffic which often dominates a job's completion time. This traffic is particularly pronounced in iterative distributed operations where each iteration invokes a shuffle operation. We observe that data of different iterations are related according to the transformation logic of distributed operations. If data generated by the current iteration are partitioned to the computers where they will be processed in the next iteration, unnecessary shuffle network traffic between the two iterations can be prevented. We model general iterative distributed operations as the transform-and-shuffle primitive and define a powerful notion named Confluence key dependency to precisely capture the data relations in the primitive. We further find that by binding key partitions between different iterations based on the Confluence key dependency, the shuffle network traffic can always be reduced by a predictable percentage. We implemented the Confluence system. Confluence provides a simple interface for programmers to express the Confluence key dependency, based on which Confluence automatically generates efficient key partitioning schemes. Evaluation results on diverse real-life applications show that Confluence greatly reduces the shuffle network traffic, resulting in as much as 23 percent job completion time reduction.
Feng Liang 0004, Francis C. M. Lau 0001, Heming Cui, Cho-Li Wang
IEEE Trans. Parallel Distributed Syst.4
2017 Scalable Adaptive NUMA-Aware Lock
abstract
Scalable locking is a key building block for scalable multi-threaded software. Its performance is especially critical in multi-socket, multi-core machines with non-uniform memory access (NUMA). Previous schemes such as in-place locks and delegation locks only perform well under a certain level of contention, and often require non-trivial tuning for a particular configuration. Besides, in large NUMA systems, current delegation locks cannot perform satisfactorily due to lack of optimized NUMA policies. In this work, we propose SANL, a locking scheme that can deliver high performance under various contention levels by adaptively switching between in-place locks and delegation locks. To optimize the performance of delegation locks, we introduce a new NUMA policy that jointly considers node distances and server utilization when choosing lock servers. We have implemented SANL and evaluated it with four popular multi-threaded applications (Memcached, Berkeley DB, Phoenix2 and SPLASH-2), on a 40-core Intel machine and a 64-core AMD machine. The comparison results with seven other representative locking schemes show that SANL outperforms them in most contention situations. For example, in one group test, SANL is 3.7 times faster than RCL lock and 17 times faster than POSIX mutex.
Haibo Chen 0001, Luwei Cheng, Francis C. M. Lau 0001, Cho-Li Wang
IEEE Trans. Parallel Distributed Syst.5
2016 Lightweight Dependency Checking for Parallelizing Loops with Non-Deterministic Dependency on GPU
abstract
General-purpose GPUs have been prevalent for a decade. Nevertheless, GPU programming remains an onerous job practically exclusive to veteran developers who must know both domain-specific knowledge and GPU architecture well. Although current parallelizing compilers that automatically parallelize and offload sizable loops onto the GPU have helped in unfettering the power of the GPU with minimal programming effort, there are still a family of loops that carry statically non-deterministic data dependencies and cannot be parallelized. To tackle this issue, we propose two lightweight dependency checking schemes that are very different from existing conservative compilers to assist parallelizing loops with non-deterministic data dependencies. Our schemes feature linear work complexity for memory operations, lower memory consumption compared to previous work, and minimal false positives by leveraging the lockstep execution on the GPU's SIMD lanes. Experiments done using microbenchmarking and real-life applications on the latest advanced AMD discrete GPUs show that our schemes can achieve 2.2 × speedup over existing solutions in dependency-free cases while only taking about 20% of time compared to existing solutions in the case with statically unproven loop-carried dependencies.
Hongyuan Liu 0002, King Tin Lam, Huanxin Lin, Cho-Li Wang
ICPADS4
2016 Scalable adaptive NUMA-aware lock: combining local locking and remote locking for efficient concurrency
abstract
Scalable locking is a key building block for scalable multi-threaded software. Its performance is especially critical in multi-socket, multi-core machines with non-uniform memory access (NUMA). Previous schemes such as local locking and remote locking only perform well under a certain level of contention, and often require non-trivial tuning for a particular configuration. Besides, for large NUMA systems, because of unmanaged lock server's nomination, current distance-first NUMA policies cannot perform satisfactorily.
Francis C. M. Lau 0001, Cho-Li Wang, Luwei Cheng, Haibo Chen 0001
PPoPP3
2015 Cache Affinity Optimization Techniques for Scaling Software Transactional Memory Systems on Multi-CMP Architectures
abstract
Software transactional memory (STM) enhances both ease-of-use and concurrency, and is considered one of the next-generation paradigms for parallel programming. Application programs may see hotspots where data conflicts are intensive and seriously degrade the performance. So advanced STM systems employ dynamic concurrency control techniques to curb the conflict rate through properly throttling the rate of spawning transactions. High-end computers may have two or more multicore processors so that data sharing among cores goes through a non-uniform cache memory hierarchy. This poses challenges to concurrency control designs as improper metadata placement and sharing will introduce scalability issues to the system. Poor thread-to-core mappings that induce excessive cache invalidation are also detrimental to the overall performance. In this paper, we share our experience in designing and implementing a new dynamic concurrency controller for Tiny STM, which helps keeping the system concurrency at a near-optimal level. By decoupling unfavourable metadata sharing, our controller design avoids costly inter-processor communications. It also features an affinity-aware thread migration technique that fine-tunes thread placements by observing inter-thread transactional conflicts. We evaluate our implementation using the STAMP benchmark suite and show that the controller can bring around 21% average speedup over the baseline execution.
Kinson Chan, King Tin Lam, Cho-Li Wang
ISPDC3
2015 Cloud, grid, P2P and internet computing: Recent trends and future directions
Sang-Soo Yeo, Yu Chen 0002, Cho-Li Wang
Peer-to-Peer Netw. Appl.3
2015 Optimization of Composite Cloud Service Processing with Virtual Machines
abstract
By leveraging virtual machine (VM) technology, we optimize cloud system performance based on refined resource allocation, in processing user requests with composite services. Our contribution is three-fold. (1) We devise a VM resource allocation scheme with a minimized processing overhead for task execution. (2) We comprehensively investigate the best-suited task scheduling policy with different design parameters. (3) We also explore the best-suited resource sharing scheme with adjusteddivisible resource fractions on running tasks in terms of Proportional-share model (PSM), which can be split into absolute mode (called AAPSM) and relative mode (RAPSM). We implement a prototype system over a cluster environment deployed with 56 real VM instances, and summarized valuable experience from our evaluation. As the system runs in short supply, lightest workload first (LWF) is mostly recommended because it can minimize the overall response extension ratio (RER) for both sequential-mode tasks and parallel-mode tasks. In a competitive situation with over-commitment of resources, the best one is combining LWF with both AAPSM and RAPSM. It outperforms other solutions in the competitive situation, by 16 + % w.r.t. the worst-case response time and by 7.4 + % w.r.t. the fairness.
Sheng Di, Derrick Kondo, Cho-Li Wang
IEEE Trans. Computers3
2015 Latency-aware DVFS for efficient power state transitions on many-core architectures
Zhiquan Lai, King Tin Lam, Cho-Li Wang, Jinshu Su
J. Supercomput.3
2014 Adaptive Live VM Migration over a WAN: Modeling and Implementation
abstract
Recent advances in virtualization technology enable high mobility of virtual machines (VMs) and resource provisioning at a data-center level. Various strategies have been proposed for fast VM live migration over a local-area network (LAN). The most common solution uses memory pre-copying and assumes storage is shared on the LAN. When applied to a wide-area network (WAN), a new design philosophy in VM live migration algorithms is necessary to address like challenges of long latency, limited or unstable bandwidth and storage relocation. This paper proposes a three-phase, fractional, hybrid pre-copy and post-copy solution for both memory and storage to achieve highly adaptive and responsive WAN-wide migration. Our strategy is to selectively migrate an important fraction of memory and storage in the pre-copy and freeze-and-copy phases, while the rest (non-critical data set) is post-copied or demand-paged. We propose a new metric called performance restoration agility, which considers both the downtime and VM speed degradation during the post-copy phase, to evaluate the migration process. We also develop a profiling framework and a novel probabilistic prediction model to adaptively find a predictably optimal combination of the memory and storage fractions to migrate. Our solution is implemented on Xen and evaluated in an emulated WAN environment. Experimental results show that the solution achieves better adaptiveness than others for various applications over a WAN while retaining the responsiveness of post-copy algorithms.
Weida Zhang, King Tin Lam, Cho-Li Wang
IEEE CLOUD3
2014 Resource Allocation in Cloud Environment: A Model Based on Double Multi-attribute Auction Mechanism
abstract
In this paper, a resource allocation model is constructed, based on the Double Multi-Attribute Auction (DMAA) mechanism. Firstly, multiple attributes are taken into account to form the Quality Index (QI), which is used to comprehensively evaluate consumers' and providers' performance in the transactions. Secondly, a Support Vector Machine (SVM) algorithm is adopted to predict the price. Finally, the Mean-Variance Optimization (MVO) algorithm is solved to obtain the optimized resource allocation scheme. Simulation results show that the proposed model can improve the resource utilization while satisfying user needs better.
Xingwei Wang 0001, Cho-Li Wang, Keqin Li 0001, Min Huang 0001
CloudCom3
2014 Rhymes: A shared virtual memory system for non-coherent tiled many-core architectures
abstract
The rising core count per processor is pushing chip complexity to a level that hardware-based cache coherency protocols become too hard and costly to scale someday. We need new designs of many-core hardware and software other than traditional technologies to keep up with the ever-increasing scalability demands. A cluster-on-chip architecture, as exemplified by the Intel Single-chip Cloud Computer (SCC), promotes a software-oriented approach instead of hardware support to implementing shared memory coherence. This paper presents a shared virtual memory (SVM) system, dubbed Rhymes, tailored to new processor kinds of non-coherent and hybrid memory architectures. Rhymes features a two-way cache coherence protocol to enforce release consistency for pages allocated in shared physical memory (SPM) and scope consistency for pages in percore private memory. It also supports page remapping on a percore basis to boost data locality. We implement and test Rhymes on the SCC port of the Barrelfish OS. Experimental results show that our SVM outperforms the pure SPM approach used by Intel's software managed coherence (SMC) library by up to 12 times through improved cache utilization for applications with strong data reuse patterns.
King Tin Lam, Jinghao Shi, Dominic Hung, Cho-Li Wang, Zhiquan Lai, Wangbin Zhu, Youliang Yan
ICPADS4
2014 Adaptive Algorithm for Minimizing Cloud Task Length with Prediction Errors
abstract
Compared to traditional distributed computing like grid system, it is non-trivial to optimize cloud task's execution performance due to its more constraints like user payment budget and divisible resource demand. In this paper, we analyze in-depth our proposed optimal algorithm minimizing task execution length with divisible resources and payment budget: 1) We derive the upper bound of cloud task length, by taking into account both workload prediction errors and hostload prediction errors. With such state-of-the-art bounds, the worst-case task execution performance is predictable, which can improve the quality of service in turn. 2) We design a dynamic version for the algorithm to adapt to the load dynamics over task execution progress, further improving the resource utilization. 3) We rigorously build a cloud prototype over a real cluster environment with 56 virtual machines, and evaluate our algorithm with different levels of resource contention. Cloud users in our cloud system are able to compose various tasks based on off-the-shelf web services. Experiments show that task execution lengths under our algorithm are always close to their theoretical optimal values, even in a competitive situation with limited available resources. We also observe a high level of fair treatment on the resource allocation among all tasks.
Sheng Di, Cho-Li Wang, Franck Cappello
IEEE Trans. Cloud Comput.2
2014 Computational awareness towards green environments
Neil Y. Yen, Cho-Li Wang, Jong Hyuk Park 0001
J. Supercomput.2
2013 Towards Payment-Bound Analysis in Cloud Systems with Task-Prediction Errors
abstract
In modern cloud systems, how to optimize user service level based on virtual resources customized on demand is a critical issue. In this paper, we comprehensively analyze the payment bound under a cloud model with virtual machines (VMs), by taking into account that task's workload may be predicted with errors. The analysis is based on an optimized resource allocation algorithm with polynomial time complexity. We theoretically derive the upper bound of task payment based on a particular margin of workload prediction-error. We also extend the payment-minimization algorithm to adapt to the dynamic changes of host availability over time, and perform the evaluation by a real-cluster environment with 56 VMs deployed. Experiments confirm the correctness of our theoretical inference, and show that our payment-minimization solution can keep 95% of user payments below 1.15 times as large as the theoretical values of the ideal payment with hypothetically accurate information. The ratio for the rest user payments can be limited to about 1.5 at the worst case.
Sheng Di, Cho-Li Wang, Derrick Kondo
IEEE CLOUD2
2013 GPU-TLS: An Efficient Runtime for Speculative Loop Parallelization on GPUs
abstract
Recently GPUs have risen as one important parallel platform for general purpose applications, both in HPC and cloud environments. Due to the special execution model, developing programs for GPUs is difficult even with the recent introduction of high-level languages like CUDA and OpenCL. To ease the programming efforts, some research has proposed automatically generating parallel GPU codes by complex compile-time techniques. However, this approach can only parallelize loops 100% free of inter-iteration dependencies (i.e., DOALL loops). To exploit runtime parallelism, which cannot be proven by static analysis, in this work, we propose GPU-TLS, a runtime system to speculatively parallelize possibly-parallel loops in sequential programs on GPUs. GPU-TLS parallelizes a possibly-parallel loop by chopping it into smaller sub-loops, each of which is executed in parallel by a GPU kernel, speculating that no inter-iteration dependencies exist. After dependency checking, the buffered writes of iterations without mis-speculations are copied to the master memory while iterations encountering mis-speculations are re-executed. GPU-TLS addresses several key problems of speculative loop parallelization on GPUs: (1) The larger mis-speculation rate caused by larger number of threads is reduced by three approaches: the loop chopping parallelization approach, the deferred memory update scheme and intra-warp value forwarding method. (2) The larger overhead of dependency checking is reduced by a hybrid scheme: eager intra-warp dependency checking combined with lazy inter-warp dependency checking. (3) The bottleneck of serial commit is alleviated by a parallel commit scheme, which allows different iterations to enter the commit phase out of order but still guarantees sequential semantics. Extensive evaluations using both micro benchmarks and real-life applications on two recent NVIDIA GPU cards show that speculative loop parallelization using GPU-TLS can achieve speedups ranging from 5 to 160 for sequential programs with possibly-parallel loops.
Chenggang Zhang, Cho-Li Wang
CCGRID3
2013 Minimization of cloud task execution length with workload prediction errors
abstract
In cloud systems, it is non-trivial to optimize task's execution performance under user's affordable budget, especially with possible workload prediction errors. Based on an optimal algorithm that can minimize cloud task's execution length with predicted workload and budget, we theoretically derive the upper bound of the task execution length by taking into account the possible workload prediction errors. With such a state-of-the-art bound, the worst-case performance of a task execution with a certain workload prediction errors is predictable. On the other hand, we build a close-to-practice cloud prototype over a real cluster environment deployed with 56 virtual machines, and evaluate our solution with different resource contention degrees. Experiments show that task execution lengths under our solution with estimates of worst-case performance are close to their theoretical ideal values, in both non-competitive situation with adequate resources and the competitive situation with a certain limited available resources. We also observe a fair treatment on the resource allocation among all tasks.
Sheng Di, Cho-Li Wang
HiPC2
2013 PVTCP: Towards practical and effective congestion control in virtualized datacenters
abstract
While modern datacenters are increasingly adopting virtual machines (VMs) to provide elastic cloud services, they still rely on traditional TCP for congestion control. In virtualized datacenters, TCP endpoints are separated by a virtualization layer and subject to the intervention of the hypervisor's scheduling. Most previous attempts focused on tuning the hypervisor layer to try to improve the VMs' I/O performance, and there is very little work on how a VM's guest OS may help the transport layer to adapt to the virtualized environment. In this paper, we find that VM scheduling delays can heavily contaminate RTTs as sensed by VM senders, preventing TCP from correctly learning the physical network condition. After giving an account of the source of the problem, we propose PVTCP, a ParaVirtualized TCP to counter the distorted congestion information caused by VM scheduling on the sender side. PVTCP is self-contained, requiring no modification to the hypervisor. Experiments show that PVTCP is much more effective in addressing incast congestion in virtualized datacenters than standard TCP.
Luwei Cheng, Cho-Li Wang, Francis C. M. Lau 0001
ICNP2
2013 Java with Auto-parallelization on Graphics Coprocessing Architecture
abstract
GPU-based many-core accelerators have gained a footing in supercomputing. Their widespread adoption yet hinges on better parallelization and load scheduling techniques to utilize the hybrid system of CPU and GPU cores easily and efficiently. This paper introduces a new user-friendly compiler framework and runtime system, dubbed Japonica, to help Java applications harness the full power of a heterogeneous system. Japonica unveils an all-round system design unifying the programming style and language for transparent use of both CPU and GPU resources, automatically parallelizing all kinds of loops and scheduling workloads efficiently across the CPU-GPU border. By means of simple user annotations, sequential Java source code will be analyzed, translated and compiled into a dual executable consisting of CUDA kernels and multiple Java threads running on GPU and CPU cores respectively. Annotated loops will be automatically split into loop chunks (or tasks) being scheduled to execute on all available GPU/CPU cores. Implementing a GPU-tailored thread-level speculation (TLS) model, Japonica supports speculative execution of loops with moderate dependency densities and privatization of loops having only false dependencies on the GPU side. Our scheduler also supports task stealing and task sharing algorithms that allow swift load redistribution across GPU and CPU. Experimental results show that Japonica, on average, can run 10x, 2.5x and 2.14x faster than the best serial (1-thread CPU), GPU-alone and CPU-alone versions respectively.
Chenggang Zhang, King Tin Lam, Cho-Li Wang
ICPP4
2013 Optimization and stabilization of composite service processing in a cloud system
abstract
With virtual machines (VM), we design a cloud system aiming to optimize the overall performance, in processing user requests made up of composite services. We address three contributions. (1) We optimize VM resource allocation with a minimized processing overhead subject to task's payment budget. (2) For maximizing the fairness of treatment in a competitive situation, we investigate the best-suited scheduling policy. (3) We devise a resource sharing scheme adjusted based on Proportional-Share model, further mitigating the resource contention. Experiments confirm two points: (1) mean task response time approaches the theoretically optimal value in non-competitive situation; (2) as system runs in short supply, each request could still be processed efficiently as compared to their ideal results. Combining Lightest Workload First (LWF) policy with Adjusted Proportional-Share Model (LWF+APSM) exhibits the best performance. It outperforms others in a competitive situation, by 38% w.r.t. worst-case response time and by 12% w.r.t. fairness of treatment.
Sheng Di, Derrick Kondo, Cho-Li Wang
IWQoS3
2013 Optimization of cloud task processing with checkpoint-restart mechanism
abstract
In this paper, we aim at optimizing fault-tolerance techniques based on a checkpointing/restart mechanism, in the context of cloud computing. Our contribution is three-fold. (1) We derive a fresh formula to compute the optimal number of checkpoints for cloud jobs with varied distributions of failure events. Our analysis is not only generic with no assumption on failure probability distribution, but also attractively simple to apply in practice. (2) We design an adaptive algorithm to optimize the impact of checkpointing regarding various costs like checkpointing/restart overhead. (3) We evaluate our optimized solution in a real cluster environment with hundreds of virtual machines and Berkeley Lab Checkpoint/Restart tool. Task failure events are emulated via a production trace produced on a large-scale Google data center. Experiments confirm that our solution is fairly suitable for Google systems. Our optimized formula outperforms Young's formula by 3-10 percent, reducing wall-clock lengths by 50-100 seconds per job on average.
Sheng Di, Yves Robert, Frédéric Vivien, Derrick Kondo, Cho-Li Wang, Franck Cappello
SC5
2013 Network performance isolation for latency-sensitive cloud applications
Luwei Cheng, Cho-Li Wang
Future Gener. Comput. Syst.2
2013 Dynamic Optimization of Multiattribute Resource Allocation in Self-Organizing Clouds
abstract
By leveraging virtual machine (VM) technology which provides performance and fault isolation, cloud resources can be provisioned on demand in a fine grained, multiplexed manner rather than in monolithic pieces. By integrating volunteer computing into cloud architectures, we envision a gigantic self-organizing cloud (SOC) being formed to reap the huge potential of untapped commodity computing power over the Internet. Toward this new architecture where each participant may autonomously act as both resource consumer and provider, we propose a fully distributed, VM-multiplexing resource allocation scheme to manage decentralized resources. Our approach not only achieves maximized resource utilization using the proportional share model (PSM), but also delivers provably and adaptively optimal execution efficiency. We also design a novel multiattribute range query protocol for locating qualified nodes. Contrary to existing solutions which often generate bulky messages per request, our protocol produces only one lightweight query message per task on the Content Addressable Network (CAN). It works effectively to find for each task its qualified resources under a randomized policy that mitigates the contention among requesters. We show the SOC with our optimized algorithms can make an improvement by 15-60 percent in system throughput than a P2P Grid model. Our solution also exhibits fairly high adaptability in a dynamic node-churning environment.
Sheng Di, Cho-Li Wang
IEEE Trans. Parallel Distributed Syst.2
2013 Error-Tolerant Resource Allocation and Payment Minimization for Cloud System
abstract
With virtual machine (VM) technology being increasingly mature, compute resources in cloud systems can be partitioned in fine granularity and allocated on demand. We make three contributions in this paper: 1) We formulate a deadline-driven resource allocation problem based on the cloud environment facilitated with VM resource isolation technology, and also propose a novel solution with polynomial time, which could minimize users' payment in terms of their expected deadlines. 2) By analyzing the upper bound of task execution length based on the possibly inaccurate workload prediction, we further propose an error-tolerant method to guarantee task's completion within its deadline. 3) We validate its effectiveness over a real VM-facilitated cluster environment under different levels of competition. In our experiment, by tuning algorithmic input deadline based on our derived bound, task execution length can always be limited within its deadline in the sufficient-supply situation; the mean execution length still keeps 70 percent as high as user-specified deadline under the severe competition. Under the original-deadline-based solution, about 52.5 percent of tasks are completed within 0.95-1.0 as high as their deadlines, which still conforms to the deadline-guaranteed requirement. Only 20 percent of tasks violate deadlines, yet most (17.5 percent) are still finished within 1.05 times of deadlines.
Sheng Di, Cho-Li Wang
IEEE Trans. Parallel Distributed Syst.2
2012 Lightweight Application-Level Task Migration for Mobile Cloud Computing
abstract
Mobile cloud computing allows mobile applications to use the enormous resources in the clouds. In order to seamlessly utilize the resources, it is common to migrate computation among mobile nodes and cloud nodes. Therefore, a highly portable and transparent migration approach is needed. In terms of portability, application-level migration with code instrumentation is the most portable approach. However, in the existing literature, this approach imposes significant runtime overhead, even when no migration takes place. Most of these works are for mobile agents, and migrations are to be invoked by the programs. Migration points are also restricted to certain locations where migration status is being polled. In this paper, we propose a Java byte code transformation technique for realizing task migration without imposing significant overhead on normal execution. Asynchronous migration technique is used to allow migrations to take place virtually anywhere in the user codes, and the proposed Twin Method Hierarchy minimizes the overhead resulting from state-restoration codes in normal execution. We have implemented our approach in our middleware. The results show that our approach can allow lightweight computation migration at application level, achieve considerable speedups and utilize the cloud resources from mobile devices.
Ricky K. K. Ma, Cho-Li Wang
AINA2
2012 vBalance: using interrupt load balance to improve I/O performance for SMP virtual machines
abstract
A Symmetric MultiProcessing (SMP) virtual machine (VM) enables users to take advantage of a multiprocessor infrastructure in supporting scalable job throughput and request responsiveness. It is known that hypervisor scheduling activities can heavily degrade a VM's I/O performance, as the scheduling latencies of the virtual CPU (vCPU) eventually translates into the processing delays of the VM's I/O events. As for a UniProcessor (UP) VM, since all its interrupts are bound to the only vCPU, it completely relies on the hypervisor's help to shorten I/O processing delays, making the hypervisor increasingly complicated. Regarding SMP-VMs, most researches ignore the fact that the problem can be greatly mitigated at the level of guest OS, instead of imposing all scheduling pressure on the hypervisor.
Luwei Cheng, Cho-Li Wang
SoCC2
2012 SmartShadow-K: an practical knowledge network for joint context inference in everyday life
abstract
Smart environments require to percept conditions of people. Current context-aware systems mainly model limited user situations, which constrains their coverage and effect in real world usage. This paper proposes an encyclopedic knowledge network to enable practical context inference in our daily life by: 1) expressing essential semantics of contextual concepts and relations into a well-informed relational network, and 2) exploiting relational semantics to infer various contexts simultaneously. The performance of the approach is validated in real challenging problems and compared with inference of human being.
Li Zhang 0045, Gang Pan 0001, Zhaohui Wu 0001, Shijian Li, Cho-Li Wang
UbiComp5
2012 Mobile Edutainment with Interactive Augmented Reality Using Adaptive Marker Tracking
abstract
Augmented Reality (AR) is a great partner of Edutainment in motivating students and enriching a class. Students can check out additional digital information of physical items such as books, samples, exhibits or even sites. The information looks just like existing in the reality. In addition to a show of multimedia content, we opine that if the virtual objects augmented to the reality can act and react like real objects, the learning experience can be greatly promoted. We developed an AR system in that the virtual objects can interact with each others. We implemented several applications using the system to demonstrate the use of interactive AR in edutainment. We also attempted to resolve the limitation of image angle coming with vision-based AR by building a multi-marker mechanism using a cube structure with surface area approximation. We also discussed some challenges and issues that researchers or developers should take note of in pursuing AR development.
Chi-Lam Lai, Cho-Li Wang
ICPADS2
2012 Decentralized proactive resource allocation for maximizing throughput of P2P Grid
Sheng Di, Cho-Li Wang
J. Parallel Distributed Comput.2
2011 Towards Context-Aware Ubiquitous Transaction Processing: A Model and Algorithm
abstract
Transaction management for mobile and ubiquitous computing aims at providing mobile users with reliable services in a transparent way anytime anywhere. To make such a vision a reality, transaction processing for the mobile and ubiquitous computing needs to adapt to the runtime environments dynamically. However, most existing mobile transaction models do not consider the context-based transaction management. In this paper, we propose a context-aware transaction model and context-driven coordination algorithms. They are built on an event-context-action mechanism, enabling the transaction processing to adapt well to dynamically changing transaction context. The simulation results have also demonstrated that our model and algorithms can significantly improve the successful commit ratio under unstable context conditions.
Feilong Tang 0001, Song Guo 0001, Minyi Guo, Minglu Li 0001, Cho-Li Wang
ICC5
2011 TrC-MC: Decentralized Software Transactional Memory for Multi-multicore Computers
abstract
To achieve single-lock atomicity in software transactional memory systems, the commit procedure often goes through a common clock variable. When there are frequent transactional commits, clock sharing becomes inefficient. Tremendous cache contention takes place between the processors and the computing throughput no longer scales with processor count. Therefore, traditional transactional memories are unable to accelerate applications with frequent commits regardless of thread count. While systems with decentralized data structures have better performance on these applications, we argue they are incomplete as they create much more aborts than traditional transactional systems. In this paper we apply two design changes, namely zone partitioning and timestamp extension, to optimize an existing decentralized algorithm. We prove the correctness and evaluate some benchmark programs with frequent transactional commits. We find it as much as several times faster than the state-of-the-art software transactional memory system. We have also reduced the abort rate of the system to an acceptable level.
Kinson Chan, Cho-Li Wang
ICPADS2
2011 Probabilistic Best-Fit Multi-dimensional Range Query in Self-Organizing Cloud
abstract
With virtual machine (VM) technology being increasingly mature, computing resources in modern Cloud systems can be partitioned in fine granularity and allocated on demand with "pay-as-you-go" model. In this work, we study the resource query and allocation problems in a Self-Organizing Cloud (SOC), where host machines are connected by a peer-to-peer (P2P) overlay network on the Internet. To run a user task in SOC, the requester needs to perform a multi-dimensional range search over the P2P network for locating host machines that satisfy its minimal demand on each type of resources. The multi-dimensional range search problem is known to be challenging as contentions along multiple dimensions could happen in the presence of the uncoordinated analogous queries. Moreover, low resource matching rate may happen while restricting query delay and network traffic. We design a novel resource discovery protocol, namely Proactive Index Diffusion CAN (PID-CAN), which can proactively diffuse resource indexes over the nodes and randomly route query messages among them. Such a protocol is especially suitable for the range query that needs to maximize its best-fit resource shares under possible competition along multiple resource dimensions. Via simulation, we show that PID-CAN could keep stable and optimized searching performance with low query delay and traffic overhead, for various test cases under different distributions of query ranges and competition degrees. It also performs satisfactorily in dynamic node-churning situation.
Sheng Di, Cho-Li Wang, Weida Zhang, Luwei Cheng
ICPP2
2011 WAVNet: Wide-Area Network Virtualization Technique for Virtual Private Cloud
abstract
A Virtual Private Cloud (VPC) is a secure collection of computing, storage and network resources spanning multiple sites over Wide Area Network (WAN). With VPC, computation and services are no longer restricted to a fixed site but can be relocated dynamically across geographical sites to improve manageability, performance and fault tolerance. We propose WAVNet, a layer 2 virtual private network (VPN) which supports virtual machine live migration over WAN to realize mobility of execution environment across multiple security domains. WAVNet adopts a UDP hole punching technique to achieve direct network connection between two Internet hosts without special router configuration. We evaluate our design in an emulated WAN with 64 hosts and also in a real WAN environment with 10 machines located at seven different sites across the Asia-Pacific region. The experimental results show that WAVNet not only achieves close-to-native host-to-host network bandwidth and latency, but also guarantees more effective VM live migration than existing solutions.
Zheming Xu, Sheng Di, Weida Zhang, Luwei Cheng, Cho-Li Wang
ICPP5
2010 BetterLife 2.0: Large-Scale Social Intelligence Reasoning on Cloud
abstract
This paper presents the design of the Better Life 2.0 framework, which facilitates implementation of large-scale social intelligence application in cloud environment. We argued that more and more mobile social applications in pervasive computing need to be implemented this way, with a lot of user generated activities in social networking websites. We adopted the Case-based Reasoning technique to provide logical reasoning and outlined design considerations when porting a typical CBR framework jCOLIBRI2 to cloud, using Hadoop's various services (HDFS, HBase). These services allow efficient case base management (e.g. case insertion) and distribution of computational intensive jobs to speed up reasoning process more than 5 times. With the scalability merit of MapReduce, we can improve recommendation service with social network analysis that needs to handle millions of users' social activities.
Dexter H. Hu, Yinfeng Wang, Cho-Li Wang
CloudCom3
2010 Optimizing data acquisition by sensor-channel co-allocation in wireless sensor networks
abstract
Wireless sensor networks (WSNs) should handle multiple sensing tasks for various applications. How to improve the quality of the data acquired in such resource constrained environment is a challenging issue. In this paper, we propose a sensor-channel co-allocation model for scheduling the sensing tasks. The proposed model considers the capability, coupling and load balancing constraints for sensing data acquisition, and can guarantee transmission of sensed data in real-time while avoiding data incompleteness in an efficient way. A spatiotemporal metric called sensing-span is proposed to evaluate the tasks' execution cost of achieving desired data quality. We extend computation task scheduling a lgorithms to support sensor-channel co-allocation problem and a heuristic called Minimum Service Capability Fragment (MSCF) is introduced for task scheduling to minimize the waste of reserved channel capacity. Simulation results show that MSCF can improve the performance of data acquisition in WSNs as compared with other heuristics, when scheduling a large number of concurrent data acquisition tasks.
Yinfeng Wang, Cho-Li Wang, Jiannong Cao 0001, Alvin Chan Toong Shoon
HiPC2
2010 Dual-Phase Just-in-Time Workflow Scheduling in P2P Grid Systems
abstract
This paper presents a fully decentralized just-in-time workflow scheduling method in a P2P Grid system. The proposed solution allows each peer node to autonomously dispatch inter-dependent tasks of workflows to run on geographically distributed computers. To reduce the workflow completion time and enhance the overall execution efficiency, not only does each node perform as a scheduler to distribute its tasks to execution nodes (or resource nodes), but the resource nodes will also set the execution priorities for the received tasks. By taking into account the unpredictability of tasks' finish time, we devise an efficient task scheduling heuristic, namely dynamic shortest makespan first (DSMF), which could be applied at both scheduling phases for determining the priority of the workflow tasks. We compare the performance of the proposed algorithm against seven other heuristics by simulation. Our algorithm achieves 20%~60% reduction on the average completion time and 37.5%~90% improvement on the average workflow execution efficiency over other decentralized algorithms.
Sheng Di, Cho-Li Wang
ICPP2
2010 A Stack-on-Demand Execution Model for Elastic Computing
abstract
Cloud computing is all the rage these days; its confluence with mobile computing would bring an even more pervasive influence. Clouds per se are elastic computing infrastructure where mobile applications can offload or draw tasks in an on-demand push-pull manner. Lightweight and portable task migration support enabling better resource utilization and data access locality is the key for success of mobile cloud computing. Existing task migration mechanisms are however too coarse-grained and costly, offsetting the benefits from migration and hampering flexible task partitioning among the mobile and cloud resources. We propose a new computation migration technique called stack-on-demand (SOD) that exports partial execution states of a stack machine to achieve agile mobility, easing into small-capacity devices and flexible distributed execution in a multi-domain workflow style. Our design also couples SOD with a novel object faulting technique for efficient access to remote objects. We implement the SOD concept into a middleware system for transparent execution migration of Java programs. It is shown that SOD migration cost is pretty low, comparing to several existing migration mechanisms. We also conduct experiments with an iPhone handset to demonstrate the elasticity of SOD by which server-side heavyweight processes can run adaptively on the cell phone.
Ricky K. K. Ma, King Tin Lam, Cho-Li Wang, Chenggang Zhang
ICPP3
2010 Adaptive sampling-based profiling techniques for optimizing the distributed JVM runtime
abstract
Extending the standard Java virtual machine (JVM) for cluster-awareness is a transparent approach to scaling out multithreaded Java applications. While this clustering solution is gaining momentum in recent years, efficient runtime support for fine-grained object sharing over the distributed JVM remains a challenge. The system efficiency is strongly connected to the global object sharing profile that determines the overall communication cost. Once the sharing or correlation between threads is known, access locality can be optimized by collocating highly correlated threads via dynamic thread migrations. Although correlation tracking techniques have been studied in some page-based software DSM systems, they would entail prohibitively high overheads and low accuracy when ported to fine-grained object-based systems. In this paper, we propose a lightweight sampling-based profiling technique for tracking inter-thread sharing. To preserve locality across migrations, we also propose a stack sampling mechanism for profiling the set of objects which are tightly coupled with a migrant thread. Sampling rates in both techniques can vary adaptively to strike a balance between preciseness and overhead. Such adaptive techniques are particularly useful for applications whose sharing patterns could change dynamically. The profiling results can be exploited for effective thread-to-core placement and dynamic load balancing in a distributed object sharing environment. We present the design and preliminary performance result of our distributed JVM with the profiling implemented. Experimental results show that the profiling is able to obtain over 95% accurate global sharing profiles at a cost of only a few percents of execution time increase for fine- to medium-grained applications.
King Tin Lam, Cho-Li Wang
IPDPS3
2010 GPS Calibrated Ad-Hoc Localization for Geosocial Networking
Dexter H. Hu, Cho-Li Wang, Yinfeng Wang
UIC2
2009 A Case-Based Component Selection Framework for Mobile Context-Aware Applications
abstract
This paper proposes a new semi-reliable multicast algorithm based on the (m, k)-firm scheduling technique, where in each consecutive k window messages sent by a sender, at least m of these messages must be received by the receiver. To assure this restriction, message recovery mechanisms from reliable multicast protocols can be used. This algorithm is mainly applicable to applications that may suffer losses, as long as these losses do not occur consecutively and do not overrun a specified maximum value of message, without any degradation of the quality of service.
Dexter H. Hu, Cho-Li Wang
ISPA4
2009 SmartShadow: Modeling A User-centric Mobile Virtual Space
abstract
This paper attempts to model pervasive computing environments as a user-centric ldquoSmartShadowrdquo using the BDP (belief-desire-plan) user model, which maps pervasive computing environments into a dynamic virtual user space. SmartShadow will follow the user to provide him with pervasive services, just like his shadow in the physical world. In the BDP model, desires of a user are inferred from his belief set, and plans are made to satisfy each desire. Pervasive service is introduced to describe computing resources in the cyberspace, which can be organized by the user's BDP to accomplish his desires. The composition process maps pervasive services into a user's SmartShadow. The model is logically natural and simple, and can flexibly model dynamics of pervasive computing spaces. In addition, we implement a simulation system to verify and evaluate the SmartShadow model.
Li Zhang 0045, Gang Pan 0001, Zhaohui Wu 0001, Shijian Li, Cho-Li Wang
PerCom5
2008 A Performance Study of Clustering Web Application Servers with Distributed JVM
abstract
A Distributed Java Virtual Machine (DJVM) is a cluster-wide set of extended JVMs that enables parallel execution of a multithreaded Java application. It has proven effectiveness for scaling scientific applications. However, leveraging DJVMs to cluster real-life web applications with commercial server workloads has not been well studied. This paper presents a new generic clustering approach based on DJVMs that promote user transparency and global object sharing for web application servers. We port Apache Tomcat to our JESSICA2 DJVM and study the performance of a wide range of web applications running on the server. Our experimental results show that this approach can scale better than the traditional clustering approach, particularly for cache-centric web applications.
King Tin Lam, Cho-Li Wang
ICPADS3
2008 Process reassignment with reduced migration cost in grid load rebalancing
abstract
We study the load rebalancing problem in a heterogeneous grid environment that supports process migration. Given an initial assignment of tasks to machines, the problem consists of finding a process reassignment that achieves a desired better level of load balance with minimum reassignment (process migration) cost. Most previous algorithms for related problems aim mainly at improving the balance level (or makespan) with no explicit concern for the reassignment cost. We propose a heuristic which is based on local search and several optimizing techniques which include the guided local search strategy and the multi-level local search. The searching integrates both the change of workload and the migration cost introduced by a process movement into the movement selection, and enables a good tradeoff between low-cost movements and the improving balance level. Evaluations show that the proposed heuristic can find a solution with much lower migration cost for achieving the same balance level than previous greedy or local search algorithms for a range of problem cases.
Lin Chen 0025, Cho-Li Wang, Francis C. M. Lau 0001
IPDPS2
2008 Scalable group-based checkpoint/restart for large-scale message-passing systems
abstract
The ever increasing number of processors used in parallel computers is making fault tolerance support in large-scale parallel systems more and more important. We discuss the inadequacies of existing system-level checkpointing solutions for message-passing applications as the system scales up. We analyze the coordination cost and blocking behavior of two current MPI implementations with checkpointing support. A group-based solution combining coordinated checkpointing and message logging is then proposed. Experiment results demonstrate its better performance and scalability than LAM/MPI and MPICH-VCL. To assist group formation, a method to analyze the communication behaviors of the application is proposed.
Justin C. Y. Ho, Cho-Li Wang, Francis C. M. Lau 0001
IPDPS2
2008 Lightweight process migration and memory prefetching in openMosix
abstract
We propose a lightweight process migration mechanism and an adaptive memory prefetching scheme called AMPoM (adaptive memory prefetching in openMosix), whose goal is to reduce the migration freeze time in openMosix while ensuring the execution efficiency of migrants. To minimize the freeze time, our system transfers only a few pages to the destination node during process migration. After the migration, AMPoM analyzes the spatial locality of memory access and iteratively prefetches memory pages from remote to hide the latency of inter-node page faults. AMPoM adopts a unique algorithm to decide which and how many pages to prefetch. It tends to prefetch more aggressively when a sequential access pattern is developed, when the paging rate of the process is high or when the network is busy. This advanced strategy makes AMPoM highly adaptive to different application behaviors and system dynamics. The HPC Challenge benchmark results show that AMPoM can avoid 98% of migration freeze time while preventing 85-99% of page fault requests after the migration. Compared to openMosix which does not have remote page fault, AMPoM induces a modest overhead of 0-5% additional runtime. When the working set of a migrant is small, AMPoM outperforms openMosix considerably due to the reduced amount of data transfer. These results indicate that by exploiting memory access locality and prefetching, process migration can be a lightweight operation with little software overhead in remote paging.
Roy S. C. Ho, Cho-Li Wang, Francis C. M. Lau 0001
IPDPS2
2008 Object co-location and memory reuse for Java programs
abstract
We introduce a new memory management system, STEMA, which can improve the execution time of Java programs. STEMA detects prolific types on-the-fly and co-locates their objects in a special memory space which supports reuse of memory. We argue and show that memory reuse and co-location of prolific objects can result in improved cache locality, reduced memory fragmentation, reduced GC time, and faster object allocation. We evaluate STEMA using 16 benchmarks. Experimental results show that STEMA performs 2.7%, 4.0%, and 8.2% on average better than MarkSweep, CopyMS, and SemiSpace.
Zoe C. H. Yu, Francis C. M. Lau 0001, Cho-Li Wang
ACM Trans. Archit. Code Optim.3
2007 GPS-Based Location Extraction and Presence Management for Mobile Instant Messenger
Dexter H. Hu, Cho-Li Wang
EUC2
2007 A Receiver-Coordinated Approach for Throughput Aggregation in High Bandwidth Multicast
abstract
In application-level high bandwidth multicast (HBM), physical links can be shared by multiple long-lived unicast flows. We identify several data transfer patterns which can cause suboptimal bandwidth usage of narrow links and which have not been clearly identified in previous solutions for application-level HBM. We propose a distributed solution to avoid these problematic patterns, with which end systems are coordinated and each is responsible to forward a bounded amount of data. Consequently, the outgoing traffic of each end system is balanced and limited. It avoids congestion due to merging unicast flows, which increases the utilization of the narrow links. Receivers that are close by topologically request their data in a disjoint and coordinated fashion, which leads to much reduced duplicated data at the narrow links. Simulation results show that our solution can achieve higher throughputs at the receivers, which is due to more efficient utilization of the narrow links' bandwidth, than mesh-based or multiple-tree approaches.
Mark C. M. Tsang, Cho-Li Wang, Ken C. K. Tsang, Francis C. M. Lau 0001
INFOCOM2
2006 An Adaptive Multipath Protocol for Efficient IP Handoff in Mobile Wireless Networks
abstract
Achieving IP handoff with a short latency and minimal packet loss is essential for mobile devices that roam across IP subnets. Many existing solutions require changes to be made to the network or transport layer, and they tend to suffer from long handoff latency in either soft or hard hand-off scenario, or both; and some are difficult to deploy in practice. We propose a new protocol, called the adaptive multipath protocol, to achieve efficient IP handoff. Based on link-layer signal strength measurements, two different schemes are used to handle soft and hard handoff respectively. Seamless IP handoff is achieved by using multiple transport layer connections on top of persistent link-layer connectivity during soft handoff. To achieve low hand-off latency during hard handoff, a set of distributed sessions repositories (SRs), which are independent of the end hosts, are employed. Simulation results clearly support our claims. In particular, the latency for hard handoff is found to be as low as 50% of that of Fast handoff.
Ken C. K. Tsang, Roy S. C. Ho, Mark C. M. Tsang, Cho-Li Wang, Francis C. M. Lau 0001
AINA (1)4
2006 Smart Instant Messenger in Pervasive Computing Environments
Chun-Fai Law, Sung-Ming Chan, Cho-Li Wang
GPC4
2006 A segment-based DSM supporting large shared object space
abstract
This paper introduces a software DSM that can extend its shared object space exceeding 4GB in a 32-bit commodity cluster environment. This is achieved through the dynamic memory mapping mechanism, with local hard disks as backing store. We introduce the new concept of segments with intelligent splitting to reduce network traffic, false sharing as well as adapt better to the shared memory access patterns. A priority-based swapping algorithm is designed to reduce disk accesses for efficient dynamic memory mapping, and maximize the use of disk space as shared object space. A new queue-based scheme is also devised for efficient and simple management of memory blocks. The proposed solutions were implemented in LOTS V.2, and it can outperform its previous version when running small applications, while the maximum shared object space is increased to one-third of the total free disk space available among all the nodes.
Benny Wang-Leung Cheung, Cho-Li Wang
IPDPS2
2006 G-PASS: an instance-oriented security infrastructure for Grid travelers
abstract
Abstract Grid computing unifies distributed resources via its support for the creation and use of virtual organizations (VOs), where a VO represents a collection of distributed resources to be accessed through predefined resource sharing and coordination policies. We consider a special type of mobile processes, named Grid travelers, which can travel across boundaries of VOs for the detection of resource availability, to negotiate for the approval of access privileges and to conduct remote execution. A new security infrastructure named G‐PASS is proposed to guarantee the validity and integrity of the travelers and the critical security knowledge they collect while traveling, especially while crossing some VOs. G‐PASS borrows the idea of passport and custom, as well as the procedures for people's travel in real life, to provide role‐based delegation mapping and access control. We demonstrate the power and feasibility of G‐PASS with a simulated mobile agent environment and a distributed ray‐tracing application running on multiple VOs. Various security overheads coming from migration decisions and actual agent or process migration are reported. G‐PASS can be installed with Grid Security Infrastructure (GSI) as the base, which makes it compatible with the existing Grid middleware. Copyright © 2006 John Wiley & Sons, Ltd.
Tianchi Ma, Lin Chen 0025, Cho-Li Wang, Francis C. M. Lau 0001
Concurr. Comput. Pract. Exp.3
2006 An architecture to support scalable distributed virtual environment systems on grid
Cho-Li Wang, Francis C. M. Lau 0001
J. Supercomput.2
2005 Smart Retrieval and Sharing of Information Resources Based on Contexts of User-Information Relationships
abstract
Information resources on their own present only the information they contain. The relationship between the resources and the users is usually neglected. By exploiting the relationship between the user and information resources in the form of context, which is established when the user accesses or acquires these resources, can help create smart mobile appliances. The metadata contained in a context is not merely data about data, but represents how the user and the information resource are related, which helps searching, provides clues to find related information resources, and facilitates information sharing with minimum manual effort. An electronic business name cards application is implemented to demonstrate the applicability of the idea.
Wai-Kwong Wing, Francis C. M. Lau 0001, Cho-Li Wang
AINA3
2004 Exploiting Java Objects Behavior for Memory Management and Optimizations
Zoe C. H. Yu, Francis C. M. Lau 0001, Cho-Li Wang
APLAS3
2004 PAT: a postmortem object access pattern analysis and visualization tool
abstract
Applying a cache coherence protocol capable of adapting to memory access patterns is a viable approach to improving the performance of software distributed shared memory. In this paper, we present an approach of postmortem memory access pattern analysis and visualization, which has been applied to our design of a global object space for a distributed Java Virtual Machine. The tool not only can enhance our understanding of the access patterns inherent in an application but can also help us to evaluate the effectiveness of an adaptive protocol used in the design of the global object space.
Weijian Fang, Cho-Li Wang, Wenzhang Zhu, Francis C. M. Lau 0001
CCGRID2
2004 LOTS: a software DSM supporting large object space
abstract
Software DSM provides good programmability for cluster computing, but its performance and limited shared memory space for large applications hinder its popularity. This paper introduces LOTS, a C++ runtime library supporting a large shared object space. With its dynamic memory mapping mechanism, LOTS can map more objects, lazily from the local disk to the virtual memory during access, leaving only a trace of control information for each object in the local process space. To our knowledge, LOTS is the first pure runtime software DSM supporting a shared object space larger than the local process space. Our testing shows that LOTS can utilize all the free hard disk space available to support hundreds of gigabytes of shared objects with a small overhead. The scope consistency memory model and a mixed coherence protocol allow LOTS to achieve better scalability with respect to problem size and cluster size.
Benny Wang-Leung Cheung, Cho-Li Wang, Francis C. M. Lau 0001
CLUSTER2
2004 A novel adaptive home migration protocol in home-based DSM
abstract
Home migration is used to tackle the home assignment problem in home-based software distributed shared memory systems. We propose an adaptive home migration protocol to optimize the single-writer pattern which occurs frequently in distributed applications. Our approach is unique in its use of a per-object threshold which is continuously adjusted to facilitate home migration decisions. This adaptive threshold is monotonously decreasing with increased likelihood that a particular object exhibits a lasting single-writer pattern. The threshold is tuned according to the feedback of previous home migration decisions at runtime. We implement this adaptive home migration protocol in a distributed Java virtual machine that supports truly parallel execution of multithreaded Java applications on clusters. The analysis and the experiments show that our home migration protocol demonstrates both the sensitivity to the lasting single-writer pattern and the robustness against the transient single-writer pattern. In the latter case, the protocol inhibits home migration in order to reduce the home redirection overhead.
Weijian Fang, Cho-Li Wang, Wenzhang Zhu, Francis C. M. Lau 0001
CLUSTER2
2004 A Collaborative and Semantic Data Management Framework for Ubiquitous Computing Environment
Weisong Chen, Cho-Li Wang, Francis C. M. Lau 0001
EUC2
2004 Ontology Mapping in Pervasive Computing Environment
C. Y. Kong, Cho-Li Wang, Francis C. M. Lau 0001
EUC2
2004 Context-Aware State Management for Ubiquitous Applications
Pauline P. L. Siu, Nalini Moti Belaramani, Cho-Li Wang, Francis C. M. Lau 0001
EUC3
2004 State-On-Demand Execution for Adaptive Component-based Mobile Agent Systems
Yuk Chow, Wenzhang Zhu, Cho-Li Wang, Francis C. M. Lau 0001
ICPADS3
2004 Petri-Net-Based Coordination Algorithms for Grid Transactions
Feilong Tang 0001, Minglu Li 0001, Joshua Zhexue Huang, Cho-Li Wang, Zongwei Luo
ISPA4
2004 Gamelet: A Mobile Service Component for Building Multi-server Distributed Virtual Environment on Grid
Cho-Li Wang, Francis C. M. Lau 0001
ISPA2
2003 Lightweight Transparent Java Thread Migration for Distributed JVM
abstract
A distributed JVM on a cluster can provide a high-performance platform for running multithreaded Java applications transparently. Efficient scheduling of Java threads among cluster nodes in a distributed JVM is desired for maintaining a balanced system workload so that the application can achieve maximum speedup. We present a transparent thread migration system that is able to support high-performance native execution of multi-threaded Java programs. To achieve migration transparency, we perform dynamic native code instrumentation inside the JIT compiler. The mechanism has been successfully implemented and integrated in JESSICA2, a JIT-enabled distributed JVM, to enable automatic thread distribution and dynamic load balancing in a cluster environment. We discuss issues related to supporting transparent Java thread migration in a JIT-enabled distributed JVM, and compare our solution with previous approaches that use static bytecode instrumentation and JVMDI. We also propose optimizations including dynamic register patching and pseudo-inlining that can reduce the runtime overhead incurred in a migration act. We use measured experimental results to show that our system is efficient and lightweight.
Wenzhang Zhu, Cho-Li Wang, Francis C. M. Lau 0001
ICPP2
2003 Functionality Adaptation: A Context-Aware Service Code Adaptation for Pervasive Computing Environments
abstract
Pervasive computing has attracted a lot of attention in recent years. There are now proxy servers that are specially designed for pervasive computing. To enable content viewing in small devices, different kinds of content adaptation techniques have been used (such as distillation and transcoding) to adapt Web contents in content-rich servers to resource-constrained devices. Adaptation of Web contents has been widely discussed, but little attention was paid to the adaptation of services (or service code), which is equally important for computing anytime, anywhere, and on any device. We present an approach to adaptation of service code which is proxy-based and context-aware, called "functionality adaptation". The main difficulty of such an adaptation is to estimate the resource usage required for an execution, which varies with the input size and is available only at run-time. We propose a conservative solution. A simple prototype has been implemented to evaluate our adaptation approach.
VivienWai-Man Kwan, Francis C. M. Lau 0001, Cho-Li Wang
Web Intelligence3
2003 p-Jigsaw: a cluster-based Web server with cooperative caching support
abstract
Abstract Clustering provides a viable approach to building scalable Web systems with increased computing power and abundant storage space for data and contents. In this paper, we present a pure‐Java‐based parallel Web server system, p‐Jigsaw, which operates on a cluster and uses the technique of cooperative caching to achieve high performance. We introduce the design of an in‐memory cache layer, called Global Object Space (GOS), for dynamic caching of frequently requested Web objects. The GOS provides a unified view of cluster‐wide memory resources for achieving location‐transparent Web object accesses. The GOS relies on cooperative caching to minimize disk accesses. A requested Web object can be fetched from a server node's local cache or a peer node's local cache, with the disk serving only as the last resort. A prototype system based on the W3C Jigsaw server has been implemented on a 16‐node PC cluster. Three cluster‐aware cache replacement algorithms were tested and evaluated. The benchmark results show good speedups with a real‐life access log, proving that cooperative caching can have significant positive impacts on the performance of cluster‐based parallel Web servers. Copyright © 2003 John Wiley & Sons, Ltd.
Cho-Li Wang, Francis C. M. Lau 0001
Concurr. Comput. Pract. Exp.2
2003 A Grid Middleware for Distributed Java Computing with MPI Binding and Process Migration Supports
Lin Chen 0025, Cho-Li Wang, Francis C. M. Lau 0001
J. Comput. Sci. Technol.2
2003 Document replication and distribution in extensible geographically distributed web servers
Ling Zhuo, Cho-Li Wang, Francis C. M. Lau 0001
J. Parallel Distributed Comput.2
2003 On the design of global object space for efficient multi-threading Java computing on clusters
Weijian Fang, Cho-Li Wang, Francis C. M. Lau 0001
Parallel Comput.2
2003 Solving irregularly structured problems based on distributed object model
Cho-Li Wang
Parallel Comput.2
2002 M-JavaMPI: A Java-MPI Binding with Process Migration Support
abstract
Several Java bindings to the Message Passing Interface (MPI) software have been developed for high-performance parallel Java-based computing with message-passing in the past. None of them however addressed the issue of supporting transparent Java process migration for achieving dynamic load distribution and balancing. This paper presents a middleware, called M-JavaMPI, that runs on top of the standard JVM to support transparent Java process migration and communication redirection. The middleware allows Java processes to freely and transparently migrate between machines to achieve load balancing, and migrated processes can continue communication with other processes using MPI. The method we use to achieve process migration is to capture execution context and restoring the execution context at the Java bytecode level using the Java Virtual Machine Debugger Interface (JVMDI). Post-migration interprocess communication is enabled via a Restorable Java-MPI API. Tests using a 16-node cluster have Shown that our mechanism yields considerable performance gain through migration.
Ricky K. K. Ma, Cho-Li Wang, Francis C. M. Lau 0001
CCGRID2
2002 Socket Cloning for Cluster-Based Web Servers
abstract
Cluster-based web server is a popular solution to meet the demand of the ever-growing web traffic. However existing approaches suffer from several limitations to achieve this. Dispatcher-based systems either can achieve only coarse-grained load balancing or would introduce heavy load to the dispatcher Mechanisms like cooperative caching consume much network resources when transferring large cache objects. In this paper, we present a new network support mechanism, called Socket Cloning (SC), in which an opened socket can be migrated efficiently between cluster nodes. With SC, the processing of HTTP requests can be moved to the node that has a cached copy of the requested document, thus bypassing any object transfer between peer servers. A prototype has been implemented and tests show that SC incurs less overhead than all the mentioned approaches. In trace-driven benchmark tests, our system outperforms these approaches by more than 30% with a cluster of twelve web server nodes.
Yiu-Fai Sit, Cho-Li Wang, Francis C. M. Lau 0001
CLUSTER2
2002 JESSICA2: A Distributed Java Virtual Machine with Transparent Thread Migration Support
abstract
A distributed Java Virtual Machine (DJVM) spanning multiple cluster nodes can provide a true parallel execution environment for multi-threaded Java applications. Most existing DJVMs suffer from the slow Java execution in interpretive mode and thus may not be efficient enough for solving computation-intensive problems. We present JESSICA2, a new DJVM running in JIT compilation mode that can execute multi-threaded Java applications transparently on clusters. JESSICA2 provides a single system image (SSI) illusion to Java applications via an embedded global object space (GOS) layer. It implements a cluster-aware Java execution engine that supports transparent Java thread migration for achieving dynamic load balancing. We discuss the issues of supporting transparent Java thread migration in a JIT compilation environment and propose several lightweight solutions. An adaptive migrating-home protocol used in the implementation of the GOS is introduced. The system has been implemented on x86-based Linux clusters and significant performance improvements over the previous JESSICA system have been observed.
Wenzhang Zhu, Cho-Li Wang, Francis C. M. Lau 0001
CLUSTER2
2002 Efficient Global Object Space Support for Distributed JVM on Cluster
abstract
We present the design of a global object space in a distributed Java Virtual Machine that supports parallel execution of a multi-threaded Java program on a cluster of computers. The global object space virtualizes a single Java object heap across machine boundaries to facilitate transparent object accesses. Based on the object connectivity information that is available at runtime, the object reachable from threads at different nodes, called a distributed-shared object, are detected With the detection of distributed-shared objects, we can alleviate overheads in maintaining the memory consistency within the global object space. Several runtime optimization methods have been incorporated in the global object space design, including an object home migration method that reallocates the home of a distributed-shared object, synchronized method migration that allows the remote execution of a synchronized method at the home node of its synchronized object, and object pushing that uses the object connectivity information to improve access locality.
Weijian Fang, Cho-Li Wang, Francis C. M. Lau 0001
ICPP2
2002 Load Balancing in Distributed Web Server Systems with Partial Document Replication
abstract
How documents of a Web site are replicated and where they are placed among the server nodes have an important bearing on balance of load in a geographically distributed Web server (DWS) system. The traffic generated due to movements of documents at runtime could also affect the performance of the DWS system. In this paper, we prove that minimizing such traffic is NP-hard. We propose a new document distribution scheme that periodically performs partial replication of a site's documents at selected server locations to maintain load balancing. Several approximation algorithms are used in it to minimize traffic generated. The simulation results show that this scheme can achieve better load balancing than a dynamic scheme, while the internal traffic it causes has a negligible effect on the system's performance.
Ling Zhuo, Cho-Li Wang, Francis C. M. Lau 0001
ICPP2
2002 A New Asynchronous Parallel Evolutionary Algorithm for Function Optimization
Pu Liu, Francis C. M. Lau 0001, Michael J. Lewis, Cho-Li Wang
PPSN4
2002 Directed Point: a communication subsystem for commodity supercomputing with Gigabit Ethernet
Cho-Li Wang, Anthony T. C. Tam, Benny Wang-Leung Cheung, Wenzhang Zhu, David C. M. Lee 0002
Future Gener. Comput. Syst.1
2002 Portable and Scalable Algorithm for Irregular All-to-All Communication
Wenheng Liu, Cho-Li Wang, Viktor Prasanna 0001
J. Parallel Distributed Comput.2
2002 Special section on Industrial information systems: progresses and perspectives in Pacific Rim
Choon Seong Leem, Cho-Li Wang
J. Syst. Softw.2
2001 Document Distribution Algorithm for Load Balancing on an Extensible Web Server Architecture
abstract
Access latency and load balancing are the two main issues in the design of clustered Web server architecture for achieving high performance. We propose a novel document distribution algorithm for load balancing on a cluster of distributed Web servers. We group Web pages that are likely to be accessed during a request session into a migrating unit, which is used as the basic unit of document placement. A modified binning algorithm is developed to distribute the migrating units among the Web servers to fulfil the load balancing. We also present a redirection mechanism, which makes use of a migrating unit's property, to reduce the cost of request redirections. The distribution of Web documents would be recomputed periodically to adapt to the changes in client request patterns and system configuration. Simulation results show that our solution can reduce the amount of request redirection and document migration, and it can distribute workload properly among Web servers.
Ben Chung-Pun Ng, Cho-Li Wang
CCGRID2
2001 Building a Scalable Web Server with Global Object Space Support on Heterogeneous Clusters
abstract
Clustering provides a viable approach to building a scalable Web server system. Many existing cluster-based Web servers, however, do not fully utilize the underlying features of the cluster environment, and most parallel web servers are designed for homogeneous clusters. In this paper, we present a pure-Java-implemented parallel Web server that can run on heterogeneous clusters. The core of the proposed system is an application-level “global object space”, which is an integration of the available physical memory of the cluster nodes for storing frequently requested objects. The global object space provides a unified view of cluster-wide memory resources, and allows transparent accesses to cached objects. Using a technique known as cooperative caching, a requested Web object can be fetched from a node’s local memory cache or a peer node’s memory cache to avoid hot spots and excessive disk operations. A preliminary prototype system has been implemented by modifying the W3C’s Jigsaw Web server. We obtained good speedups in the benchmark tests, indicating that clustering with cooperative caching can greatly improve the performance of a Web server system. 1.
Cho-Li Wang, Francis C. M. Lau 0001
CLUSTER2
2001 A Distributed Object Model for Solving Irregularly Structured Problems on Cluster
abstract
This paper presents a distributed object model MOIDE (Multithreaded Object-oriented Infrastructure on Distributed Environment) for solving irregularly structured problems on cluster. The primary appeal of MOIDE is its flexible system structure that is adaptive to heterogeneous architecture of a cluster. MOIDE integrates the object-oriented and multithreaded methodologies to set up a unified computing environment. Both the shared-data access and remote messaging are incorporated in a two-layer communication mechanism for efficient inter-object communication with the common communication interface. MOIDE supports dynamic load balancing by its autonomous load scheduling technique. A runtime support system implements the MOIDE model as a platform-independent infrastructure for developing and executing irregularly structured applications. N-body, ray tracing, and conjugate gradient applications are implemented to illustrate the advantages of MOIDE model.
Cho-Li Wang
CLUSTER2
2001 Distributed particle simulation method on adaptive collaborative system
Zhengyu Liang, Cho-Li Wang
Future Gener. Comput. Syst.3
2000 Contention-free Complete Exchange Algorithm on Clusters
abstract
To construct a large commodity cluster a hierarchical network is generally adopted for connecting the host machines, where a Gigabit backbone switch connects a few commodity switches with uplinks to achieve scaled bisectional bandwidth. This type of interconnection usually results in link contention and has congestion developed at the uplink ports. Moreover the non-deterministic delays on scheduling communication events in clusters accelerate the building up of congestion amongst these uplink ports, which lead to severe packets drop and hinder the overall performance. In this paper, we focus on the practical design of high-speed complete exchange algorithm on a commodity cluster interconnected by a hierarchical Ethernet-based network. By exploiting some architectural characteristics of the interconnection in optimizing the performance of a complete exchange algorithm, we introduce a congestion control mechanism-global windowing that monitors and regulates the traffic load, together with a permutation scheme-reorder scheme that effectively alleviates the congestion problem. We evaluate our algorithm and compare its performance with other algorithms in a PC cluster connected by various types of switches, including Gigabit Ethernet, input-buffered and shared-memory fast Ethernet switches.
Anthony T. C. Tam, Cho-Li Wang
CLUSTER2
2000 JESSICA: Java-Enabled Single-System-Image Computing Architecture
Matchy J. M. Ma, Cho-Li Wang, Francis C. M. Lau 0001
J. Parallel Distributed Comput.2
1999 Push-Pull Messaging: A High-Performance Communication Mechanism for Commodity SMP Clusters
abstract
Push-Pull Messaging is a novel messaging mechanism for high-speed interprocess communication in a cluster of symmetric multi-processors (SMP) machines. This messaging mechanism exploits the parallelism in SMP nodes by allowing the execution of communication stages of a messaging event on different processors to achieve maximum performance. Push-Pull Messaging facilitates further improvement on communication performance by employing three optimizing techniques in our design: (1) Cross-Space Zero Buffer provides a unified buffer management mechanism to achieve a copy-less communication for the data transfer among processes within a SMP node. (2) Address Translation Overhead Masking removes the address translation overhead from the critical path in the internode communication. (3) Push-and-Acknowledge Overlapping overlaps the push and acknowledge phases to hide the acknowledge latency. Overall, Push-Pull Messaging effectively utilizes the system resources and improves the communication speed. It has been implemented to support high-speed communication for connecting quad Pentium Pro SMPs with 100 Mbit/s Fast Ethernet.
Kwan-Po Wong, Cho-Li Wang
ICPP2
1999 Resource Scaling Effects on MPP Performance: The STAP Benchmark Implications
abstract
Presently, massively parallel processors (MPPs) are available only in a few commercial models. A sequence of three ASCI Teraflops MPPs has appeared before the new millenium. This paper evaluates six MPP systems through STAP benchmark experiments. The STAP is a radar signal processing benchmark which exploits regularly structured SPMD data parallelism. We reveal the resource scaling effects on MPP performance along orthogonal dimensions of machine size, processor speed, memory capacity messaging latency, and network bandwidth. We show how to achieve balanced resources scaling against enlarged workload (problem size). Among three commercial MPPs, the IBM SP2 shows the highest speed and efficiency, attributed to its well-designed network with middleware support for single system image. The Cray T3D demonstrates a high network bandwidth with a good NUMA memory hierarchy. The Intel Paragon trails far behind due to slow processors used and excessive latency experienced in passing messages. Our analysis projects the lowest STAP speed on the ASCI Red, compared with the projected speed of two ASCI Blue machines. This is attributed to slow processors used in ASCI Red and the mismatch between its hardware and software. The Blue Pacific shows the highest potential to deliver scalable performance up to thousands of nodes. The Blue Mountain is designed to have the highest network bandwidth. Our results suggest a limit on the scalability of the distributed shared-memory (DSM) architecture adopted in Blue Mountain. The scaling model offers a quantitative method to match resource scaling with problem scaling to yield a truly scalable performance. The model helps MPP designers optimize the processors, memory, network, and I/O subsystems of an MPP. For MPP users, the scaling results can be applied to partition a large workload for SPMD execution or to minimize the software overhead in collective communication or remote memory update operations. Finally, our scaling model is assessed to evaluate MPPs with benchmarks other than STAP.
Kai Hwang 0001, Choming Wang, Cho-Li Wang
IEEE Trans. Parallel Distributed Syst.3
1998 Parallel Algorithms for Perceptual Grouping on Distributed Memory Machines
Yongwha Chung, Cho-Li Wang, Viktor Prasanna 0001
J. Parallel Distributed Comput.2
1997 Evaluating MPI Collective Communication on the SP2, T3D, and Paragon Multicomputers
abstract
We evaluate the architectural support of collective communication operations on the IBM SP2, Cray T3D, and Intel Paragon. The MPI performance data are obtained from the STAP benchmark experiments jointly performed at the USC and HKU. The T3D demonstrated clearly the best timing performance in almost all collective operations. This is attributed to the special hardware built in the T3D for fast messaging and block data transfer. With hardwired barriers, the T3D performs the barrier synchronization in 3 /spl mu/s at least 30 times faster than the SP2 or Paragon. The startup latency of collective operations increases either linearly or logarithmically in three multicomputers. For short messages, the SP2 outperforms the Paragon in the barrier, total exchange, scatter, and gather operations. Various collective operations with 64 KBytes per message over 64 nodes of the three machines can be completed in the time range (5.12 ms, 675 ms). The Paragon outperforms the SP2 in almost all collective operations with long messages. We have derived closed-form expressions to quantify the collective messaging times and aggregated bandwidth on all three machines. For total exchange with 64 nodes, the T3D, Paragon, and SP2 achieved an aggregated bandwidth of 1.745, 0.879, and 0.818 GBytes/s, respectively. These findings are useful to those who wish to predict the MPP performance or to optimize parallel applications by trade-offs between divided computation and collective communication.
Kai Hwang 0001, Choming Wang, Cho-Li Wang
HPCA3
1996 Portable Message Passing Algorithms for Irregular All-to-all Communication
abstract
In this paper we develop portable and scalable algorithms for performing irregular all-to-all communication in High Performance Computing (HPC) systems. To minimize the communication latency, the algorithm reduces the total number of messages transmitted, reduces the variance of the lengths of these messages, and overlaps the communication with computation. The performance of the algorithm is characterized using a simple model of HPC systems. Our implementations are performed using the Message Passing Interface (MPI) standard and they can be ported to various HPC platforms. The performance of our algorithms is evaluated on CM5, T3D and SP2. The results show the effectiveness of the techniques as well as the interplay between the architectural features, the machine size, and the variance of message lengths. The experiences of our study can be applied in other HPC systems to optimize the performance of collective communication operations.
W. H. Liu, Cho-Li Wang, Viktor Prasanna 0001
ICDCS2
1996 High-performance computing for vision
abstract
The main focus of the paper is on effectively using commercial-off-the-shelf (COTS) based general purpose parallel computing platforms to realize high speed implementations of vision tasks. Due to the successful use of the COTS-based systems in a variety of high performance applications, it is attractive to consider their use for vision applications as well. However, the irregular data dependencies in vision tasks lead to large communication overheads in the HPC systems. At the University of Southern California, our research efforts have been directed toward designing scalable parallel algorithms for vision tasks on the HPC systems. In our approach, we use the message passing programming model to develop portable code. Our algorithms are specified using C and MPI. In this paper, we summarize our efforts, and illustrate our approach using several example vision tasks.
Cho-Li Wang, Prashanth B. Bhat, Viktor Prasanna 0001
Proc. IEEE1
1994 Scalable parallel implementations of perceptual grouping on connection machine CM-5
abstract
Perceptual grouping is a key step in vision to organize image data into structural hypotheses to be used for high level analysis. We propose data allocation and load balancing strategies which reduce the communication cost and evenly distribute the grouping operations among the processors. These techniques result in scalable algorithms for performing perceptual grouping on CM-5. The performance of our algorithms depends only on the total grouping operations generated by the image data and is independent of the distribution of the data among the processors. Our implementations show that given a 1 K/spl times/1 K input image, extraction of line segments and several perceptual grouping steps can be performed in 5.0 seconds using a partition of CM-5 having 32 processing nodes. A serial implementation of these steps on a Sun Sparc 400 takes more than 2 minutes.
Viktor Prasanna 0001, Cho-Li Wang
ICPR (3)2
1994 Scalable Data Parallel Implementations of Object Recognition Using Geometric Hashing
Cho-Li Wang, Viktor Prasanna 0001, Hyoung Joong Kim, Ashfaq Khokhar 0001
J. Parallel Distributed Comput.1
1992 An architecture for tree search based vector quantization for single chip implementation
abstract
Vector quantization (VQ) has become feasible for use in real-time applications by employing VLSI technology. The authors propose a new search algorithm and an architecture for implementing it, which can be used in real-time image processing. This search algorithm takes O(k) time units on a sequential machine, where k is the dimension of the codevectors, assuming unit time corresponds to one comparison operation. The proposed architecture employs a single processing element (PE) and O(N) external memory for storing N hyperplanes used in the search, where N is the number of codevectors. Compared with known architectures for VQ in the literature, the proposed design does not perform any multiplication operation, since the search method is independent of any L/sub q/ metric, 1>
Heonchul Park, Viktor Prasanna 0001, Cho-Li Wang
ASAP3