VLDB 2026 Research / reviewers in the wild / expert
Dan Huang 0001
dblp:41/3836-1
· DBLP profile ↗
50ranked-venue papers
9as first author
37since 2021 · last 2026
0000-0001-5582-1031ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 42 · 9 first-author · 29 since 2021Applied, interdisciplinary, general and emerging computing · 4 · 4 since 2021Artificial intelligence and machine learning · 2 · 2 since 2021Computer networks · 2 · 2 since 2021Software engineering, systems software and programming languages · 2 · 2 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | FLARE: Fine-Grained Length-Aware Routing for Resource-Efficient Heterogeneous LLM ServingabstractWith the rapid proliferation of large language models (LLMs), model pools have become increasingly heterogeneous in both capability and efficiency.Larger LLMs can improve quality but incur higher latency and cost, while smaller LLMs are the opposite, making perquery model selection crucial in practice.This has spawned LLM routers that dispatch each query to an appropriate model.Existing routers lack fine-grained resource awareness across deployment settings, which degrades efficiency metrics in real-world serving.To this end, we propose FLARE, a length-centric, resourceaware multi-LLM routing framework that uses length-based models to estimate per-query latency and cost.FLARE formulates routing as a discrete multi-objective optimization problem to achieve an efficient trade-off.Experiments show that FLARE reduces latency and cost by up to 68% and 75% while achieving sufficient accuracy, and can be easily applied to new datasets and LLMs. Yujia Fu, Heming Zhong, Dan Huang 0001, Yutong Lu |
ACL (1) | 3 |
| 2026 | PolyKAN: A High-Performance and Universal GPU Operator Library for Polynomial Kolmogorov-Arnold NetworksabstractKolmogorov–Arnold Networks (KANs) promise higher expressive capability and stronger interpretability than Multilayer Perceptron, particularly in the domain of AI for Science. However, practical adoption has been hindered by low GPU utilization of existing parallel implementations. To address this challenge, we present a GPU-accelerated operator library, named PolyKAN, which is the first general open-source implementation of KAN and its variants. PolyKAN fuses the forward and backward passes of polynomial KAN layers into a concise set of optimized CUDA kernels. Four orthogonal techniques underpin the design: (i) lookup-table with linear interpolation that replaces runtime expensive math-library functions; (ii) 2D tiling to expose thread-level parallelism with preserving memory locality; (iii) a two-stage reduction scheme converting scattered atomic updates into a single controllable merge step; and (iv) coefficient-layout reordering yielding unit-stride reads under the tiled schedule. PolyKAN can deliver 1.2–10 × faster inference and 1.4–12 × faster training than a Triton + cuBLAS baseline, with identical accuracy on speech, audio-enhancement, and tabular-regression workloads on both highend GPU and consumer-grade GPU. Mingkun Yu, Heming Zhong, Jiazhi Jiang, Dan Huang 0001, Yutong Lu |
ICS | 4 |
| 2026 | ASM-SpMM: Unleashing the Potential of Arm SME for Sparse Matrix Multiplication AccelerationabstractSparse Matrix–Matrix Multiplication (SpMM) is a core kernel in scientific computing, data analytics, and artificial intelligence, supporting applications such as linear solvers and Graph Neural Networks (GNNs). The Scalable Matrix Extension (SME) in Armv9 introduces dedicated matrix acceleration for ARM CPUs, but exploiting its full potential for SpMM requires architecture-aware optimizations to address irregular sparsity and hardware constraints. Jiazhi Jiang, Xijia Yao, Jinhui Wei, Dan Huang 0001, Yutong Lu |
PPoPP | 5 |
| 2026 | Efficient KV Cache Spillover Management on Memory-Constrained GPU for LLM InferenceabstractThe rapid growth of model parameters presents a significant challenge when deploying large generative models on GPU. Existing LLM runtime memory management solutions tend to maximize batch size to saturate GPU device utilization. Nevertheless, this practice leads to situations where the KV Cache of certain sequences cannot be accommodated on GPUs with limited memory capacity during the model inference, requiring temporary eviction from GPU memory (referred to as KV Cache spillover). However, without careful consideration of the LLM inference's runtime pattern, current LLM inference memory management solutions face issues like one-size-fits-all spillover handling approach for different platforms, under-utilization of GPU in prefill stage, and suboptimal sequence selection due to direct employment of swap or recomputation. In this paper, we introduce FuseSpill, a holistic KV Cache management solution designed to boost LLM inference on memory-constrained GPU by efficiently handling KV Cache spillover. Specifically, FuseSpill consists of a spillover cost model that analyzes the system cost of spillover handling techniques quantitatively, a KV cache swap orchestrator to further refine the basic swap technique to sophisticated disaggregate KV Cache across heterogeneous devices for decoding iterations, a multi-executor scheduler to effectively coordinate task executors across devices, and a response length predictor to exploit the length-aware sequence selection strategy when KV Cache spillover occurs. The experimental results demonstrate that our implementation outperforms existing solutions, delivering a 20% to 40% increase in throughput while simultaneously reducing the inference latency of the spillover sequences. Jiazhi Jiang, Yao Chen 0008, Zining Zhang 0001, Bingsheng He, Pingyi Luo, Mian Lu, Yuqiang Chen, Hongbin Zhang 0006, Jiangsu Du, Dan Huang 0001, Yutong Lu |
IEEE Trans. Parallel Distributed Syst. | 10 |
| 2025 | Doppeladler: Adaptive Tensor Parallelism for Latency-Critical LLM Deployment on CPU-GPU Integrated End-User DeviceabstractLLM deployment on end-user devices has attracted significant interest from tech giants and research institutions due to privacy benefits and the elimination of network roundtrips. Reducing latency is crucial for improving user experience. Enduser devices, such as desktop and mobile processors, often integrate CPU and GPU on a single die, making tensor parallelism a promising approach to distribute workloads and reduce inference latency. However, the predefined tensor parallelism in traditional practices cannot guarantee optimization due to heterogeneity and resource contention in the CPU-GPU integrated end-user devices. In this paper, we propose Doppeladler, a practical framework designed to facilitate parallel inference of LLM on end-user devices. Doppeladler adaptively optimizes tensor parallelism based on real-time conditions and device status, enhancing resource utilization by refining workload partitioning and scheduling with a heuristic-based approach. Additionally, it lowers the costs of CPU-GPU hybrid execution by managing device resource, trigger workload rebalance during runtime to mitigate the penalty of contention, and minimizing communication overhead through zero-copy techniques and alternate access pattern. Experiments demonstrate that Doppeladler outperforms existing methods by 1.2 to 2.5 times. Jiazhi Jiang, Jiangsu Du, Dan Huang 0001, Yutong Lu |
PACT | 4 |
| 2025 | IasRT: Interference-Aware and SLO-Driven GPU Scheduling for Real-Time DNN InferenceabstractDeep Neural Network (DNN) inference has become a cornerstone of latency-sensitive applications such as autonomous driving and augmented reality. While GPUs offer high throughput for DNN inference, they often suffer from underutilization due to coarse-grained scheduling and limited concurrency. Existing GPU-sharing methods either lack awareness of fine-grained kernel interference or fail to meet service-level objectives (SLOs) under multi-tenant, multi-priority workloads. In this paper, we propose IasRT, a runtime framework that enables interferenceaware and SLO-driven GPU sharing for real-time DNN inference. IasRT profiles kernel-level resource usage and interference sensitivity, and dynamically partitions GPU streaming multiprocessors (SMs) to collocate jobs with minimal performance degradation. Furthermore, it introduces a dynamic SLO controller to maintain latency targets for multiple latency-sensitive (LS) jobs simultaneously. Evaluations on real-world DNN workloads show that IasRT reduces the 99th percentile latency of LS jobs by up to 38% compared to the state-of-the-art GPU sharing methods, while maintaining similar overall throughput from multiple collocated workloads, demonstrating its effectiveness in high-concurrency environments. Heming Zhong, Jinhui Wei, Yujia Fu, Dan Huang 0001, Yutong Lu |
ICCD | 4 |
| 2025 | HStencil: Matrix-Vector Stencil Computation with Interleaved Outer Product and MLAabstractStencil computations are fundamental to various HPC and intelligent computing applications, often consuming significant execution time. The emergence of specialized matrix units presents new opportunities to accelerate stencil computations. While scalable matrix compute units provide substantial computing horsepower, prior efforts fail to fully utilize the computing capabilities for stencils due to suboptimal matrix-unit utilization, limited instruction-level parallelism, and low cache hit rates. This paper introduces HStencil, a novel stencil computing framework utilizing matrix and vector units. HStencil addresses these challenges through three contributions: 1) microkernels that jointly leverage matrix and vector units to enhance hardware utilization; 2) fine-grained instruction scheduling with interleaved execution to enhance instruction-level parallelism; and 3) spatial prefetch to sustain high performance when working sets exceed cache capacity. Evaluations on representative benchmarks demonstrate that HStencil achieves maximum speedups of 1.81x – 5.76x over auto-vectorization across different CPU platforms, delivers 31% - 91% higher performance versus state-of-the-art methods. Jiabin Xie, Guangnan Feng, Xianwei Zhang 0001, Dan Huang 0001, Zhiguang Chen 0001, Yutong Lu |
SC | 5 |
| 2025 | GPU acceleration for DNA sequence alignment algorithm and its application
Heming Zhong, Xiaojian Pan, Zengquang He, Haoling Wang, Dan Huang 0001, Zhiguang Chen 0001 |
CCF Trans. High Perform. Comput. | 5 |
| 2025 | Co-Designing Transformer Architectures for Distributed Inference With Low CommunicationabstractTransformer models have shown significant success in a wide range of tasks. However, the massive resources required for its inference prevent deployment on a single device with relatively constrainted resources, thus leaving a high threshold of integrating their advancements. Observing scenarios such as smart home applications on edge devices and cloud deployment on commodity hardware, it is promising to distribute Transformer inference across multiple devices. Unfortunately, due to the tightly-coupled feature of Transformer model, existing model parallelism approaches necessitate frequent communication to resolve data dependencies, making them unacceptable for distributed inference, especially under relatively weak interconnection. In this paper, we propose DeTransformer, a communication-efficient distributed Transformer inference system. The key idea of DeTransformer involves the co-design of Transformer architecture to reduce the communication during distributed inference. In detail, DeTransformer is based on a novel block parallelism approach, which restructures the original Transformer layer with a single block to the decoupled layer with multiple sub-blocks. Thus, it can exploit model parallelism between sub-blocks. Next, DeTransformer contains an adaptive execution approach that strikes a trade-off among communication capability, computing power and memory budget over multiple devices. It incorporates a two-phase planning for execution, namely static planning and runtime planning. The static planning runs offline, containing a profiling procedure and a weight placement strategy before execution. The runtime planning dynamically determines the optimal parallel computing strategy from an expertly crafted search space based on real-time requests. Notably, this execution approach can adapt to heterogeneous devices by distributing workload based on devices’ computing capabilities. We conduct experiments for both auto-regressive and auto-encoder tasks of Transformer models. Experimental results show that DeTransformer can reduce distributed inference latency by up to 2.81× compared to the SOTA approach on 4 devices, while effectively maintaining task accuracy and a consistent model size. Jiangsu Du, Yuanxin Wei, Shengyuan Ye, Jiazhi Jiang, Xu Chen 0004, Dan Huang 0001, Yutong Lu |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2025 | Critique of "Productivity, Portability, Performance Data-Centric Python" by SCC Team From Sun Yat-sen UniversityabstractIn SC21, Ziogas et al. proposed Data-Centric (DaCe) Python. It attains high performance and portability, and further extends the original productivity of Python. This paper analyzes the reproducibility of the DaCe paper as part of the SC22 Student Cluster Competition (SCC). The reproduction experiments are conducted on the Azure CycleCloud. Different from the DaCe paper, we use AMD EPYC 7V73X processors for CPU-based experiments. We successfully reproduce most of the results of the DaCe paper. The remaining results are also explainable. Tengyang Zheng, Tianxing Yang, Siran Liu, Shengyou Lu, Guangnan Feng, Zhiguang Chen 0001, Dan Huang 0001 |
IEEE Trans. Parallel Distributed Syst. | 10 |
| 2024 | Communication-Efficient Model Parallelism for Distributed In-Situ Transformer InferenceabstractTransformer models have shown significant success in a wide range of tasks. Meanwhile, massive resources required by its inference prevent scenarios with resource-constrained devices from in-situ deployment, leaving a high threshold of integrating its advances. Observing that these scenarios, e.g. smart home of edge computing, are usually comprise a rich set of trusted devices with untapped resources, it is promising to distribute Transformer inference onto multiple devices. However, due to the tightly-coupled feature of Transformer model, existing model parallelism approaches necessitate frequent communication to resolve data dependencies, making them unacceptable for distributed inference, especially under weak interconnect of edge scenarios. In this paper, we propose DeTransformer, a communication-efficient distributed in-situ Transformer inference system for edge scenarios. DeTransformer is based on a novel block parallelism approach, with the key idea of restructuring the original Trans-former layer with a single block to the decoupled layer with multi-ple sub-blocks and exploit model parallelism between sub-blocks. Next, DeTransformer contains an adaptive placement approach to automatically select the optimal placement strategy by striking a trade-off among communication capability, computing power and memory budget. Experimental results show that DeTransformer can reduce distributed inference latency by up to 2.81 x compared to the SOTA approach on 4 devices, while effectively maintaining task accuracy and a consistent model size. Yuanxin Wei, Shengyuan Ye, Jiazhi Jiang, Xu Chen 0004, Dan Huang 0001, Jiangsu Du, Yutong Lu |
DATE | 5 |
| 2024 | Efficient Coupling Streaming AI and Ensemble Simulations on HPC Clusters
Jiazhi Jiang, Hongbin Zhang 0006, Deyin Liu, Jiangsu Du, Xiaojiao Yao, Jinhui Wei, Pin Chen, Dan Huang 0001, Yutong Lu |
Euro-Par (1) | 8 |
| 2024 | Equivariant Diffusion for Crystal Structure PredictionabstractIn addressing the challenge of Crystal Structure Prediction (CSP), symmetry-aware deep learning models, particularly diffusion models, have been extensively studied, which treat CSP as a conditional generation task. However, ensuring permutation, rotation, and periodic translation equivariance during diffusion process remains incompletely addressed. In this work, we propose EquiCSP, a novel equivariant diffusion-based generative model. We not only address the overlooked issue of lattice permutation equivariance in existing models, but also develop a unique noising algorithm that rigorously maintains periodic translation equivariance throughout both training and inference processes. Our experiments indicate that EquiCSP significantly surpasses existing models in terms of generating accurate structures and demonstrates faster convergence during the training process. Peijia Lin, Pin Chen, Qing Mo, Jianhuan Cen, Wenbing Huang 0001, Yang Liu 0005, Dan Huang 0001, Yutong Lu |
ICML | 8 |
| 2024 | Understanding the Inference Performance of Spatial Temporal Diffusion Transformer
Yuanxin Wei, Jiangsu Du, Dan Huang 0001, Nong Xiao 0001 |
NPC (1) | 4 |
| 2024 | Liger: Interleaving Intra- and Inter-Operator Parallelism for Distributed Large Model InferenceabstractDistributed large model inference is still in a dilemma where balancing cost and effect. The online scenarios demand intraoperator parallelism to achieve low latency and intensive communications makes it costly. Conversely, the inter-operator parallelism can achieve high throughput with much fewer communications, but it fails to enhance the effectiveness. Jiangsu Du, Jinhui Wei, Jiazhi Jiang, Shenggan Cheng, Dan Huang 0001, Zhiguang Chen 0001, Yutong Lu |
PPoPP | 5 |
| 2024 | APTMoE: Affinity-Aware Pipeline Tuning for MoE Models on Bandwidth-Constrained GPU NodesabstractRecently, the sparsely-gated Mixture-Of-Experts (MoE) architecture has garnered significant attention. To benefit a wider audience, fine-tuning MoE models on more affordable clusters, which are typically a limited number of bandwidthconstrained GPU nodes, holds promise. However, it is non-trivial to apply existing cost-effective fine-tuning approaches to MoE models, due to the increased ratio of data to computation. In this paper, we introduce APTMoE, which employs affinityaware pipeline parallelism for fine-tuning MoE models on bandwidth-constrained GPU nodes. We propose an affinity-aware offloading technique that enhances pipeline parallelism for both computational efficiency and model size, and it benefits from a hierarchical loading strategy and a demand-priority scheduling strategy. To improve the computation efficiency and reduce the data movement volume, the hierarchical loading strategy designs three loading phases and efficiently allocates computation across GPUs and CPUs during these phases, leveraging different levels of expert popularity and computation affinity. With the aim of alleviating the mutual interference among the three loading phases and maximizing the bandwidth utilization, the demand-priority scheduling strategy proactively and dynamically coordinates the loading execution order. Experiments demonstrate that APTMoE outperforms existing methods in most cases. Particularly, APTMoE successfully fine-tunes a 61.2B MoE model on 4 Nvidia A800 GPUs(40GB) and achieves up to $33 \%$ throughput improvement compared to the SOTA method. Yuanxin Wei, Jiangsu Du, Jiazhi Jiang, Xianwei Zhang 0001, Dan Huang 0001, Nong Xiao 0001, Yutong Lu |
SC | 6 |
| 2024 | HTDcr: a job execution framework for high-throughput computing on supercomputers
Jiazhi Jiang, Dan Huang 0001, Yutong Lu, Xiangke Liao |
Sci. China Inf. Sci. | 2 |
| 2024 | SAIH: A Scalable Evaluation Methodology for Understanding AI Performance Trend on HPC Systems
Jiangsu Du, Yingpeng Wen, Jiazhi Jiang, Dan Huang 0001, Xiangke Liao, Yutong Lu |
J. Comput. Sci. Technol. | 5 |
| 2024 | Topo: Towards a fine-grained topological data processing framework on Tianhe-3 supercomputer
Yutong Lu, Zhuo Tang, Dan Huang 0001, Zhiguang Chen 0001 |
J. Parallel Distributed Comput. | 5 |
| 2024 | AdaNAS: Adaptively Postprocessing With Self-Supervised Neural Architecture Search for Ensemble Rainfall ForecastsabstractPrevious post-processing studies on rainfall forecasts using numerical weather prediction (NWP) mainly focus on statistics-based aspects, while learning-based aspects are rarely investigated. Although some manually-designed models are proposed to raise accuracy, they are customized networks, which need to be repeatedly tried and verified, at a huge cost in time and labor. Therefore, a self-supervised neural architecture search (NAS) method without significant manual efforts called AdaNAS is proposed in this study to perform rainfall forecast post-processing and predict rainfall with high accuracy. In addition, we design a rainfall-aware search space to significantly improve forecasts for high-rainfall areas. Furthermore, we propose a rainfall-level regularization function to eliminate the effect of noise data during the training. Validation experiments have been performed under the cases ofNone,Light,Moderate,HeavyandViolenton a large-scale precipitation benchmark named TIGGE. Finally, the average mean-absolute error (MAE) and average root-mean-square error (RMSE) of the proposed AdaNAS model are 0.98 and 2.04 mm/day, respectively. Additionally, the proposed AdaNAS model is compared with other neural architecture search methods and previous studies. Compared results reveal the satisfactory performance and superiority of the proposed AdaNAS model in terms of precipitation amount prediction and intensity classification. Concretely, the proposed AdaNAS model outperformed previous best-performing manual methods with MAE and RMSE improving by 80.5% and 80.3%, respectively. Yingpeng Wen, Weijiang Yu, Fudan Zheng, Dan Huang 0001, Nong Xiao 0001 |
IEEE Trans. Geosci. Remote. Sens. | 4 |
| 2024 | Sophisticated Orchestrating Concurrent DLRM Training on CPU/GPU PlatformabstractRecommendation systems are essential to the operation of the majority of internet services, with Deep Learning Recommendation Models (DLRMs) serving as a crucial component. However, due to distinct computation, data access, and memory usage characteristics of recommendation models, the trainning of DLRMs may suffer from low resource utilization on prevalent heterogeneous CPU-GPU hardware platforms. Furthermore, as the majority of high-performance computing systems presently depend on multi-GPU computing nodes, the challenge of addressing low resource utilization becomes even more pronounced. Existing concurrent training solutions cannot be straightforwardly applied to DLRM due to various factors, such as insufficient fine-grained memory management and the lack of collaborative CPU-GPU scheduling. In this paper, we introduce RMixer, a scheduling framework that addresses these challenges by providing an efficient job management and scheduling mechanism for DLRM training jobs on heterogeneous CPU-GPU platforms. To facilitate training co-location, we first estimate the peak memory consumption of each job. Additionally, we track and collect resource utilization for DLRM training jobs. Based on the information of computational patterns, a batched job dispatcher with dynamic resource-complementary scheduling policy is proposed to co-locate DLRM training jobs on CPU-GPU platform. Scheduling strategies for both intra-GPU and inter-GPU scenarios were meticulously devised, with a focus on thoroughly examining individual GPU resource utilization and achieving a balanced state across multiple GPUs. Experimental results demonstrate that our implementation achieved up to 5.3× and 7.5× higher throughput on single GPU and 4 GPU respectively for training jobs involving various recommendation models. Rui Tian 0001, Jiazhi Jiang, Jiangsu Du, Dan Huang 0001, Yutong Lu |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2023 | Enhancing Multi-physics Coupling on ARM Many-Core Cluster
Wencheng Shi, Jiangsu Du, Dan Huang 0001, Yutong Lu |
APPT | 4 |
| 2023 | Accelerating Inference of 3D-CNN on ARMMany-core CPU via Hierarchical Model PartitionabstractMany applications such as biomedical analysis and scientific data analysis involve analyzing volumetric data. This spawns huge demand for 3D CNN. Although accelerators such as GPU may provide higher throughput on deep learning applications, they may not be available in all scenarios. CPU, especially many-core CPU, remains an attractive choice for deep learning in many scenarios. In this paper, we propose a inference solution that targets on the emerging ARM many-core CPU platform. A hierarchical partition approach is claimed to accelerate 3D-CNN inference by exploiting characteristics of memory and cache on ARM many-core CPU. Jiazhi Jiang, Zijiang Huang, Dan Huang 0001, Jiangsu Du, Yutong Lu |
DATE | 3 |
| 2023 | MixRec: Orchestrating Concurrent Recommendation Model Training on CPU-GPU platformabstractThe development of deep learning recommendation models (DLRM) and recommendation systems has significantly improved the precision of information matching. Due to distinct computation, data access, and memory usage characteristics of recommendation models, they may suffer from low resource utilization on prevalent heterogeneous CPU-GPU hardware platforms. Existing concurrent training solutions cannot be directly applied to DLRM due to various factors, such as insufficient fine-grained memory management and the lack of collaborative CPU-GPU scheduling. In this paper, we introduce MixRec, a scheduling framework that addresses these challenges by pro-viding an efficient job management and scheduling mechanism for DLRM training jobs on heterogeneous CPU-GPU platforms. To facilitate training co-location, we first estimate the peak memory consumption of each job. Additionally, we track and collect resource utilization for DLRM training jobs. Based on the information of resource usage, a batched job dispatcher with dynamic resource-complementary scheduling policy is proposed to co-locate DLRM training jobs on CPU-GPU platform. Experimental results demonstrate that our implementation achieved up to 4.42× higher throughput and 3.97× higher resource utilization for training jobs involving various recommendation models. Jiazhi Jiang, Rui Tian 0001, Jiangsu Du, Dan Huang 0001, Yutong Lu |
ICCD | 4 |
| 2023 | Optimizing massively parallel sparse matrix computing on ARM many-core processor
Jiazhi Jiang, Jiangsu Du, Dan Huang 0001, Yutong Lu |
Parallel Comput. | 4 |
| 2023 | Improving Computation and Memory Efficiency for Real-world Transformer Inference on GPUsabstractTransformer models have emerged as a leading approach in the field of natural language processing (NLP) and are increasingly being deployed in production environments. Graphic processing units (GPUs) have become a popular choice for the transformer deployment and often rely on the batch processing technique to ensure high hardware performance. Nonetheless, the current practice for transformer inference encounters computational and memory redundancy due to the heavy-tailed distribution of sequence lengths in NLP scenarios, resulting in low practical performance. In this article, we propose a unified solution for improving both computation and memory efficiency of the real-world transformer inference on GPUs. The solution eliminates the redundant computation and memory footprint across a transformer model. At first, a GPU-oriented computation approach is proposed to process the self-attention module in a fine-grained manner, eliminating its redundant computation. Next, the multi-layer perceptron module continues to use the word-accumulation approach to eliminate its redundant computation. Then, to better unify the fine-grained approach and the word-accumulation approach, it organizes the data layout of the self-attention module in block granularity. Since aforementioned approaches make the required memory size largely reduce and constantly fluctuate, we propose the chunk-based approach to enable a better balance between memory footprint and allocation/free efficiency. Our experimental results show that our unified solution achieves a decrease of average latency by 28% on the entire transformer model, 63.8% on the self-attention module, and reduces memory footprint of intermediate results by 7.8×, compared with prevailing frameworks. Jiangsu Du, Jiazhi Jiang, Hongbin Zhang 0006, Dan Huang 0001, Yutong Lu |
ACM Trans. Archit. Code Optim. | 5 |
| 2023 | Hierarchical Model Parallelism for Optimizing Inference on Many-core Processor via Decoupled 3D-CNN StructureabstractThe tremendous success of convolutional neural network (CNN) has made it ubiquitous in many fields of human endeavor. Many applications such as biomedical analysis and scientific data analysis involve analyzing volumetric data. This spawns huge demand for 3D-CNN. Although accelerators such as GPU may provide higher throughput on deep learning applications, they may not be available in all scenarios. CPU, especially many-core CPU with non-uniform memory access (NUMA) architecture, remains an attractive choice for deep learning inference in many scenarios. In this article, we propose a distributed inference solution for 3D-CNN that targets on the emerging ARM many-core CPU platform. A hierarchical partition approach is claimed to accelerate 3D-CNN inference by exploiting characteristics of memory and cache on ARM many-core CPU. Based on the hierarchical model partition approach, other optimization techniques such as NUMA-aware thread scheduling and optimization of 3D-img2row convolution are designed to exploit the potential of ARM many-core CPU for 3D-CNN. We evaluate our proposed inference solution with several classic 3D-CNNs: C3D, 3D-resnet34, 3D-resnet50, 3D-vgg11, and P3D. Our experimental results show that our solution can boost the performance of the 3D-CNN inference, and achieve much better scalability, with a negligible fluctuation in accuracy. When employing our 3D-CNN inference solution on ACL libraries, it can outperform naive ACL implementations by 11× to 50× on ARM many-core processor. When employing our 3D-CNN inference solution on NCNN libraries, it can outperform the naive NCNN implementations by 5.2× to 14.2× on ARM many-core processor. Jiazhi Jiang, Zijiang Huang, Dan Huang 0001, Jiangsu Du, Lin Chen 0002, Ziguang Chen, Yutong Lu |
ACM Trans. Archit. Code Optim. | 3 |
| 2023 | A Data-driven Approach to Harvesting Latent Reduced Models to Precondition Lossy Compression for Scientific DataabstractIn this paper, we propose and evaluate the idea that data need to be preconditioned prior to compression, such that they can better match the design philosophies of lossy compressors for HPC scientific data. In particular, we aim to identify a reduced model that can be utilized to transform the original data into a more compressible form. We begin with two PDE applications as a proof of concept, in which we demonstrate that a reduced model can indeed reside in the full model output, and can be utilized to improve compression ratios. A mathematical proof is also presented to show how the compression ratio is improved by the reduced model. We further explore more general dimension reduction techniques to extract the reduced model, including principal component analysis, singular value decomposition, and discrete wavelet transform. After preconditioning, the reduced model in conjunction with difference between the reduced model and full model is stored, which results in higher compression ratios. We evaluate the reduced models on ten scientific datasets, and the results show the effectiveness of our approaches. Given that there is no single method that consistently achieves the best performance, we further propose a selection strategy that guides users to select the best reduced model prior to data reduction. Huizhang Luo, Junqi Wang 0002, Zhenlu Qin, Dan Huang 0001, Qing Liu 0002, MengChu Zhou, Hong Jiang 0001 |
IEEE Trans. Big Data | 4 |
| 2023 | Full-Stack Optimizing Transformer Inference on ARM Many-Core CPUabstractThe past several years have witnessed tremendous success of transformer models in natural language processing (NLP), and their current landscape is increasingly diverse. Although GPU gradually becomes the dominating workhorse and de facto standard for deep learning, there are still many scenarios where using CPU remains a prevalent choice.Recently, ARM many-core processor starts emigrating to cloud computing and high-performance computing, which is promising to deploy transformer inference. In this paper, we identify several performance bottlenecks of existing inference runtime on many-core CPU including low-core usage, isolated thread configuration, inappropriate implementation of general matrix multiply (GEMM), and redundant computations for variable-length inputs. To tackle these problems, full-stack optimizations are conducted for these challenges from service level to operator level. We explore multi-instance parallelization at the service level to improve CPU core usage. To improve parallel efficiency of the inference runtime, we design NUMA-aware thread scheduling and a look-up table for optimal parallel configurations. The GEMM implementation is tailored for some critical modules to exploit the characteristics of transformer workload. To eliminate redundant computations, a novel storage format is designed and implemented to pack sparse data and a load balancing strategy is proposed for tasks with different sparsity. Experiments show that our implementation can outperform existing solutions by 1.1x to 6x with fixed-length inputs. For variable-length inputs, it achieves 1.9x to 8x speedups on different ARM many-core processors. Jiazhi Jiang, Jiangsu Du, Dan Huang 0001, Zhiguang Chen 0001, Yutong Lu, Xiangke Liao |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2022 | Characterizing and Optimizing Transformer Inference on ARM Many-core ProcessorabstractTransformer has experienced tremendous success and revolutionized the field of natural language processing (NLP). While GPU has become the de facto standard for deep learning computation in many cases, there are still many scenarios where using CPU for deep learning remains a prevalent choice. In particular, ARM many-core processor is emerging as a competitive candidate for HPC systems, which is promising to deploy Transformer inference. Jiazhi Jiang, Jiangsu Du, Dan Huang 0001, Yutong Lu |
ICPP | 3 |
| 2022 | Handling heavy-tailed input of transformer inference on GPUsabstractTransformer-based models achieve superior accuracy in the field of natural language processing (NLP) and start to be widely deployed in production. As a popular deployment device, graphic processing units (GPUs) basically adopt the batch processing technique for inferring transformer-based models and achieving high hardware performance. However, as the input sequence lengths of NLP tasks are generally variable and in a heavy-tailed distribution, the batch processing will bring large amounts of redundant computation and hurt the practical efficiency. Jiangsu Du, Jiazhi Jiang, Yang You 0001, Dan Huang 0001, Yutong Lu |
ICS | 4 |
| 2022 | Persistent Items Tracking in Large Data Streams Based on Adaptive SamplingabstractWe address the problem of persistent item tracking in large-scale data streams. A persistent item refers to the one that persists to occur in the stream over a long timespan. Tracking persistent items is an important and pivotal functionality for many networking and computing applications as persistent items, though not necessarily contributing significantly to the data volume, may convey valuable information on the data pattern about the stream. The state-of-the-art solutions of tracking persistent items require to know the monitoring time horizon to set the sampling rate. This limitation is further accentuated when we need to track the persistent items in recent w slots where w can be any value between 0 and T to support different monitoring granularity. Motivated by this limitation, we develop a persistent item tracking algorithm that can function without knowing the monitoring time horizon beforehand, and can thus track persistent items up to the current time t or within a certain time window at any moment. Our central technicality is adaptively reducing the sampling rate such that the total memory overhead can be limited while still meeting the target tracking accuracy. Through both theoretical and empirical analysis, we fully characterize the performance of our proposition. Lin Chen 0002, Raphael C.-W. Phan, Dan Huang 0001 |
INFOCOM | 4 |
| 2022 | Enhancing Distributed In-Situ CNN Inference in the Internet of ThingsabstractConvolutional neural networks (CNNS) enable machines to view the world as humans and become increasing prevalent for Internet of Things (IoT) applications. Instead of streaming the raw data to the cloud and executing CNN inference remotely, it would be very attractive to use local IoT devices to process as it enables IoT applications with independent decision-making ability. Since a single IoT device can hardly match the requirements of the CNN inference, especially for time-sensitive and high-accuracy tasks, the distributedin-situCNN inference becomes a potential solution. However, because of the inherently tightly coupled structure of existing CNN models, it is difficult to distribute the inference efficiently. In this article, we enhance the distributedin-situCNN inference in the IoT. We fundamentally reduce the communication overhead of distributed CNN inference by designing new loosely coupled structure (LCS). Experimental results demonstrate that LCS achieves the leading performance compared with other popular structures. Next, based on the LCS, we customize the partitioning method to reduce the synchronization points and design the decentralized asynchronous method to optimize communication in each synchronization point. To evaluate the effectiveness, we build a prototype system. When the number of IoT devices increases from 1 to 4, our system accelerates by up to$3.85\times $and reduces the memory footprint in each device by 70% with achieving a competitive accuracy and significantly outperforming other approaches. Jiangsu Du, Yunfei Du 0001, Dan Huang 0001, Yutong Lu, Xiangke Liao |
IEEE Internet Things J. | 3 |
| 2022 | Identifying challenges and opportunities of in-memory computing on large HPC systems
Dan Huang 0001, Zhenlu Qin, Qing Liu 0002, Norbert Podhorszki, Scott Klasky |
J. Parallel Distributed Comput. | 1 |
| 2022 | Optimizing small channel 3D convolution on GPU with tensor core
Jiazhi Jiang, Dan Huang 0001, Jiangsu Du, Yutong Lu, Xiangke Liao |
Parallel Comput. | 2 |
| 2021 | Optimizing Massively Parallel Winograd Convolution on ARM ProcessorabstractConvolution Neural Network (CNN) has gained a great success in deep learning applications and been accelerated by dedicated convolutional algorithms. Winograd-based algorithm can greatly reduce the number of arithmetic operations required in convolution. However, our experiments show that existing implementations in deep learning libraries cannot achieve expected parallel performance on ARM manycore CPUs with last-level cache (LLC). Compared to multicore processor, ARM manycore CPUs have more cores, more NUMA nodes and the parallel performance is more easily restricted by memory bandwidth, cache contention, NUMA configuration and etc. In this paper, we propose an optimized implementation for single-precision Winograd-based algorithm on ARM manycore CPUs. Our algorithm adjusts the data layout according to the input shape and is optimized for the characteristics of ARM processor, thus reducing the matrix transformation overhead and achieving high arithmetic intensity. We redesign the parallel algorithm for Winograd-based convolution to achieve a more efficient implementation for manycore CPUs. The experimental results with 32 cores show that for modern ConvNets, our implementation achieves speedups ranging from 3 × to 5 × over the state-of-the-art Winograd-based convolution on ARM processor. Even conducted on a set of convolutional benchmarks executing on a 128-core system with 4 NUMA nodes, the results show that our implementation can also achieve better performance than existing implementations on ARM processor. Dan Huang 0001, Zhiguang Chen 0001, Yutong Lu |
ICPP | 2 |
| 2021 | Enhancing Proportional IO Sharing on Containerized Big Data File SystemsabstractBig Data platforms recently employ resource management systems, such as YARN, Mesos, and Google Borg, to provision computational resources. These systems adopt containerization to share the computing resources in a multi-tenant setting with low performance overhead and interference. However, it may be observed that tenants often interfere with each other on the underlying Big Data File Systems (BDFS), e.g., Hadoop File System, which have been widely deployed as a persistent layer in current data centers. A solution with systematic generality is to containerize BDFS itself to isolate and allocate its IO sources to multiple tenants. To this end, we conduct analysis on the ineffectiveness of proportionally sharing BDFS IO resource via containerization. This ineffectiveness is due to the scheduler of containerization in “pseudo-starvation” status, in which most of IO requests are backlogged in BDFS rather than in containerization scheduler. Without enough backlogged IO requests, existing schedulers might have to maximize device utilization rather than enforce proportional sharing policy. To resolve this ineffectiveness issue, we develop a cross-layer system calledBDFS-Container, which containerizes BDFS at the Linux block IO level. Central to BDFS-Container, we propose and design a proactive IOPS throttling-based mechanism namedIOPS Regulator, which achieves a trade-off between maximizing IO utilization and accurately proportional IO sharing. The evaluation results show that our method can improve proportionally sharing BDFS IO resources by 74.4 percent on average. Dan Huang 0001, Jun Wang 0001, Qing Liu 0002, Nong Xiao 0001, Huafeng Wu, Jiangling Yin |
IEEE Trans. Computers | 1 |
| 2020 | A Comprehensive Study of In-Memory Computing on Large HPC SystemsabstractWith the increasing fidelity and resolution enabled by high-performance computing systems, simulation-based scientific discovery is able to model and understand microscopic physical phenomena at a level that was not possible in the past. A grand challenge that the HPC community is faced with is how to handle the large amounts of analysis data generated from simulations. In-memory computing, among others, is recognized to be a viable path forward and has experienced tremendous success in the past decade. Nevertheless, there has been a lack of a complete study and understanding of in-memory computing as a whole on HPC systems. This paper presents a comprehensive study, which goes well beyond the typical performance metrics. In particular, we assess the in-memory computing with regard to its usability, portability, robustness and internal design trade-offs, which are the key factors that of interest to domain scientists. We use two realistic scientific workflows, LAMMPS and Laplace, to conduct comprehensive studies on state-of-the-art in-memory computing libraries, including DataSpaces, DIMES, Flexpath and Decaf. We conduct cross-platform experiments at scale on two leading supercomputers, Titan at ORNL and Cori at NERSC, and summarize our key findings in this critical area. Dan Huang 0001, Zhenlu Qin, Qing Liu 0002, Norbert Podhorszki, Scott Klasky |
ICDCS | 1 |
| 2020 | Improving the efficiency of HPC data movement on container-based virtual cluster
Dan Huang 0001, Yutong Lu |
CCF Trans. High Perform. Comput. | 1 |
| 2019 | Identifying Latent Reduced Models to Precondition Lossy CompressionabstractWith the high volume and velocity of scientific data produced on high-performance computing systems, it has become increasingly critical to improve the compression performance. Leveraging the general tolerance of reduced accuracy in applications, lossy compressors can achieve much higher compression ratios with a user-prescribed error bound. However, they are still far from satisfying the reduction requirements from applications. In this paper, we propose and evaluate the idea that data need to be preconditioned prior to compression, such that they can better match the design philosophies of a compressor. In particular, we aim to identify a reduced model that can be utilized to transform the original data to a more compressible form. We begin with a case study of Heat3d as a proof of concept, in which we demonstrate that a reduced model can indeed reside in the full model output, and can be utilized to improve compression ratios. We further explore more general dimension reduction techniques to extract the reduced model, including principal component analysis, singular value decomposition, and discrete wavelet transform. After preconditioning, the reduced model in conjunction with difference between the reduced model and full model is stored, which results in higher compression ratios. We evaluate the reduced models on nine scientific datasets, and the results show the effectiveness of our approaches. Huizhang Luo, Dan Huang 0001, Qing Liu 0002, Zhenbo Qiao, Hong Jiang 0001, Jing Bi 0001, Haitao Yuan 0001, MengChu Zhou, Jinzhen Wang, Zhenlu Qin |
IPDPS | 2 |
| 2019 | Can I/O Variability Be Reduced on QoS-Less HPC Storage Systems?abstractFor a production high-performance computing (HPC) system, where storage devices are shared between multiple applications and managed in a best effort manner, I/O contention is often a major problem. In this paper, we propose a balanced messaging-based re-routing in conjunction with throttling at the middleware level. This work tackles two key challenges that have not been fully resolved in the past: whether I/O variability can be reduced on a QoS-less HPC storage system, and how to design a runtime scheduling system that can scale up to a large amount of cores. The proposed scheme uses a two-level messaging system to re-route I/O requests to a less congested storage location so that write performance is improved, while limiting the impact on read by throttling re-routing. An analytical model is derived to guide the setup of optimal throttling factor. We thoroughly analyze the virtual messaging layer overhead and explore whether the in-transit buffering is effective in managing I/O variability. Contrary to the intuition, in-transit buffer cannot completely solve the problem. It can reduce the absolute variability but not the relative variability. The proposed scheme is verified against a synthetic benchmark as well as being used by production applications. Dan Huang 0001, Qing Liu 0002, Jong Choi 0001, Norbert Podhorszki, Scott Klasky, Jeremy Logan, George Ostrouchov, Xubin He, Matthew Wolf |
IEEE Trans. Computers | 1 |
| 2019 | Harnessing Data Movement in Virtual Clusters for In-Situ ExecutionabstractAs a result of increasing data volume and velocity, Big Data science at exascale has shifted towards the in-situ paradigm, where large scale simulations run concurrently alongside data analytics. With in-situ, data generated from simulations can be processed while still in memory, thereby avoiding the slow storage bottleneck. However, running simulations and analytics together on shared resources will likely result in substantial contention if left unmanaged, as demonstrated in this work, leading to much reduced efficiency of simulations and analytics. Recently, virtualization technologies such as Linux containers have been widely applied to data centers and physical clusters to provide highly efficient and elastic resource provisioning for consolidated workloads including scientific simulations and data analytics. In this paper, we investigate to facilitate network traffic manipulation and reduce mutual interference on the network for in-situ applications in virtual clusters. In order to dynamically allocate the network bandwidth when it is needed, we adopt SARIMA-based techniques to analyze and predict MPI traffic issued from simulations. Although this can be an effective technique, the naïve usage of network virtualization can lead to performance degradation for bursty asynchronous transmissions within an MPI job. We analyze and resolve this performance degradation in virtual clusters. Dan Huang 0001, Qing Liu 0002, Scott Klasky, Jun Wang 0001, Jong Choi 0001, Jeremy Logan, Norbert Podhorszki |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2018 | Performance Evaluation and Analysis for MPI-Based Data Movement in Virtual Switch NetworkabstractVirtualization technologies have been widely deployed in data centers and private clusters to provide highly efficient and elastic resource provisioning. Further, virtualization has been extended to the network layer, known as network virtualization. For example, independent virtual switches have become the primary provider of network services for various virtual machines, such as VMware, Xen and Docker. This approach allows the physical network to be decoupled from the overlying virtual switch networks. However, network virutalization introduces performance degradation and scalability bottleneck to communication-intensive frameworks, such as MPI. We quantify and analyze the performance degradation involved with collective communications as well as bursty asynchronous transmission (BAT) in vswitch network environments. Our experiments illustrate that the performance of MPI communication can be degraded up to 5× in the virtual environment. Dan Huang 0001, Jun Wang 0001, Dezhi Han |
NAS | 1 |
| 2018 | Achieving Load Balance for Parallel Data Access on Distributed File SystemsabstractThe distributed file system, HDFS, is widely deployed as the bedrock for many parallel big data analysis. However, when running multiple parallel applications over the shared file system, the data requests from different processes/executors will unfortunately be served in a surprisingly imbalanced fashion on the distributed storage servers. These imbalanced access patterns among storage nodes are caused because a). unlike conventional parallel file system using striping policies to evenly distribute data among storage nodes, data-intensive file system such as HDFS store each data unit, referred to as chunk file, with several copies based on a relative random policy, which can result in an uneven data distribution among storage nodes; b). based on the data retrieval policy in HDFS, the more data a storage node contains, the higher probability the storage node could be selected to serve the data. Therefore, on the nodes serving multiple chunk files, the data requests from different processes/executors will compete for shared resources such as hard disk head and networkbandwidth, resulting in a degraded I/O performance. In this paper, we first conduct a complete analysis on how remote and imbalanced read/write patterns occur and how they are affected by the size of the cluster. We then propose novel methods, referred to as Opass, to optimize parallel data reads, as well as to reduce the imbalance of parallel writes on distributed file systems. Our proposed methods can benefit parallel data-intensive analysis with various parallel data access strategies. Opass adopts new matching-based algorithms to match processes to data so as to compute the maximum degree of data locality and balanced data access. Furthermore, to reduce the imbalance of parallel writes, Opass employs a heatmap for monitoring the I/O statuses of storage nodes and performs HM-LRU policy to select a local optimal storage node for serving write requests. Experiments are conducted on PRObE's Marmot 128-node cluster testbed and the results from both benchmark and well-known parallel applications show the performance benefits and scalability of Opass. Dan Huang 0001, Dezhi Han, Jun Wang 0001, Jiangling Yin, Xunchao Chen, Xuhong Zhang 0002, Jian Zhou 0004, Mao Ye 0008 |
IEEE Trans. Computers | 1 |
| 2017 | DFS-container: achieving containerized block I/O for distributed file systemsabstractToday BigData systems commonly use resource management systems such as TORQUE, Mesos, and Google Borg to share the physical resources among users or applications. Enabled by virtualization, users can run their applications on the same node with low mutual interference. Container-based virtualizations (e.g., Docker and Linux Containers) offer a lightweight virtualization layer, which promises a near-native performance and is adopted by some Big-Data resource sharing platforms such as Mesos. Nevertheless, using containers to consolidate the I/O resources of shared storage systems is still at an early stage, especially in a distributed file system (DFS) such as Hadoop File System (HDFS). To overcome this issue, we propose a distributed middleware system, DFS-Container, by further containerizing DFS. We also evaluate and analyze the unfairness of using containers to proportionally allocate the I/O resource of DFS. Based on these analyses and evaluations, we propose and implement a new mechanism, IOPS-Regulator, which improve the fairness of proportional allocation by 74.4% on average. Dan Huang 0001, Jun Wang 0001, Qing Liu 0001, Xuhong Zhang 0002, Xunchao Chen, Jian Zhou 0004 |
SoCC | 1 |
| 2017 | SideIO: A Side I/O system framework for hybrid scientific workflow
Jun Wang 0001, Dan Huang 0001, Huafeng Wu, Jiangling Yin, Xuhong Zhang 0002, Xunchao Chen |
J. Parallel Distributed Comput. | 2 |
| 2017 | Deister: A light-weight autonomous block management in data-intensive file systems using deterministic declustering distribution
Jun Wang 0001, Xuhong Zhang 0002, Junyao Zhang 0007, Jiangling Yin, Dezhi Han, Dan Huang 0001 |
J. Parallel Distributed Comput. | 7 |
| 2017 | Energy-Aware Adaptive Restore Schemes for MLC STT-RAM CacheabstractFor the sake of higher cell density while achieving near-zero standby power, recent research progress in Magnetic Tunneling Junction (MTJ) devices has leveraged Multi-Level Cell (MLC) configurations of Spin-Transfer Torque Random Access Memory (STT-RAM). However, in orderto mitigate the write disturbance in an MLC strategy, data stored in the soft bit must be restored back immediately after the hard bit switching is completed. Furthermore, as the result of MTJ feature size scaling, the soft bit can be expected to become disturbed by the read sensing current, thus requiring an immediate restore operation to ensure the data reliability. In this paper, we design and analyze a novel Adaptive Restore Scheme for Write Disturbance (ARS-WD) and Read Disturbance (ARS-RD), respectively. ARS-WD alleviates restoration overhead by intentionally overwriting soft bit lines which are less likely to be read. ARS-RD, on the other hand, aggregates the potential writes and restore the soft bit line at the time of its eviction from higher level cache. Both of these two schemes are based on a lightweight forecasting approach for the future read behavior of the cache block. Our experimental results show substantial reduction in soft bit line restore operations, delivering 17.9 percent decrease in overall energy consumption and 9.4 percent increase in IPC, while incurring negligible capacity overhead. Moreover, ARS promotes advantages of MLC to provide a preferable L2 design alternative in terms of energy, area and latency product compared to SLC STT-RAM alternatives. Xunchao Chen, Navid Khoshavi, Ronald F. DeMara, Jun Wang 0001, Dan Huang 0001, Wujie Wen, Yiran Chen 0001 |
IEEE Trans. Computers | 5 |
| 2016 | AOS: adaptive overwrite scheme for energy-efficient MLC STT-RAM cacheabstractSpin-Transfer Torque Random Access Memory (STT-RAM) has been identified as an advantageous candidate for on-chip memory technology due to its high density and ultra low leakage power. Recent research progress in Magnetic Tunneling Junction (MTJ) devices has developed Multi-Level Cell (MLC) STT-RAM to further enhance cell density. To avoid the write disturbance in MLC strategy, data stored in the soft bit must be restored back immediately after the hard bit switching is completed. However, frequent restores are not only unnecessary, but also introduce a significant energy consumption overhead. In this paper, we propose an Adaptive Overwrite Scheme (AOS) which alleviates restoration overhead by intentionally overwriting selected soft bits based on RRD (Read Reuse Distance). Our experimental results show 54.6% reduction in soft bit restoration, delivering 10.8% decrease in overall energy consumption. Moreover, AOS promotes MLC to be a preferable L2 design alternative in terms of energy, area and latency product. Xunchao Chen, Navid Khoshavi, Jian Zhou 0004, Dan Huang 0001, Ronald F. DeMara, Jun Wang 0001, Wujie Wen, Yiran Chen 0001 |
DAC | 4 |
| 2015 | Opass: Analysis and Optimization of Parallel Data Access on Distributed File SystemsabstractIn this paper, we study parallel data access on distributed file systems, e.g, the Hadoop file system. Our experiments show that parallel data read requests are often served data remotely and in an imbalanced fashion. This results in a serious disk access and data transfer contention on certain cluster/storage nodes. We conduct a complete analysis on how remote and imbalanced read patterns occur and how they are affected by the size of the cluster. We then propose a novel method to Optimize Parallel Data Access on Distributed File Systems referred to as Opass. The goal of Opass is to reduce remote parallel data accesses and achieve a higher balance of data read requests between cluster nodes. To achieve this goal, we represent the data read requests that are issued by parallel applications to cluster nodes as a graph data structure where edges weights encode the demands of data locality and load capacity. Then we propose new matching-based algorithms to match processes to data based on the configurations of the graph data structure so as to compute the maximum degree of data locality and balanced access. Our proposed method can benefit parallel data-intensive analysis with various parallel data access strategies. Experiments are conducted on PRObEs Marmot 128-node cluster tested and the results from both benchmark and well-known parallel applications show the performance benefits and scalability of Opass. Jiangling Yin, Jun Wang 0001, Jian Zhou 0004, Tyler Lukasiewicz, Dan Huang 0001, Junyao Zhang 0007 |
IPDPS | 5 |