VLDB 2026 Research / reviewers in the wild / expert
Xuanhua Shi
dblp:85/5317
· DBLP profile ↗
118ranked-venue papers
18as first author
42since 2021 · last 2026
0000-0001-8451-8656ORCID · corroborated
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 65 · 12 first-author · 17 since 2021Databases, data management, data science and information retrieval · 14 · 2 first-author · 6 since 2021Artificial intelligence and machine learning · 12 · 10 since 2021Applied, interdisciplinary, general and emerging computing · 10 · 8 since 2021Software engineering, systems software and programming languages · 8 · 1 first-author · 5 since 2021Graphics, computer vision, multimedia, augmented reality and games · 4 · 4 since 2021Security and privacy · 3 · 1 since 2021Human-computer interaction and ubiquitous computing · 2 · 1 first-authorComputer networks · 1Theory of computation · 1 · 1 first-author
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | RecToM: A Benchmark for Evaluating Machine Theory of Mind in LLM-based Conversational Recommender SystemsabstractLarge Language models (LLMs) are revolutionizing the conversational recommender systems (CRS) through their impressive capabilities in instruction comprehension, reasoning, and human interaction. A core factor underlying effective dialogue is the ability to infer and reason about others' mental states (such as desire, intention, and belief), a cognitive capacity commonly referred to as Theory of Mind (ToM). Despite growing interest in evaluating ToM in LLMs, current benchmarks predominantly rely on synthetic narratives inspired by Sally-Anne test, which emphasize physical perception and fail to capture the complexity of mental state inference in real-world conversational settings. Moreover,existing benchmarks often overlook a critical component of human ToM: behavioral prediction, the ability to use inferred mental states to guide strategic decision-making and select appropriate conversational actions for future interactions. To better align LLM-based ToM evaluation with human-like social reasoning, we propose RecToM, a novel benchmark for evaluating ToM abilities in recommendation dialogues. RecToM focuses on two complementary dimensions: Cognitive Inference and Behavioral Prediction. The former focus on understanding what has been communicated by inferring the underlying mental states. The latter emphasizes what should be done next, evaluating whether LLMs can leverage these inferred mental states to predict, select, and assess appropriate dialogue strategies. Together, these dimensions enable a comprehensive assessment of ToM reasoning in CRS. Extensive experiments on state-of-the-art LLMs demonstrate that RecToM poses a significant challenge. While the models exhibit partial competence in recognizing mental states, they struggle to maintain coherent, strategic ToM reasoning throughout dynamic recommendation dialogues, particularly in tracking evolving intentions and aligning conversational strategies with inferred mental states. Mengfan Li 0001, Xuanhua Shi, Yang Deng 0002 |
AAAI | 2 |
| 2026 | CoSToM: Causal-oriented Steering for Intrinsic Theory-of-Mind Alignment in Large Language ModelsabstractTheory of Mind (ToM), the ability to attribute mental states to others, is a hallmark of social intelligence.While large language models (LLMs) demonstrate promising performance on standard ToM benchmarks, we observe that they often fail to generalize to complex taskspecific scenarios, relying heavily on prompt scaffolding to mimic reasoning.The critical misalignment between the internal knowledge and external behavior raises a fundamental question: Do LLMs truly possess intrinsic cognition, and can they externalize this internal knowledge into stable, high-quality behaviors?To answer this, we introduce COSTOM 1 (Causal-oriented Steering for ToM alignment), a framework that transitions from mechanistic interpretation to active intervention.First, we employ causal tracing to map the internal distribution of ToM features, empirically uncovering the internal layers' characteristics in encoding fundamental ToM semantics.Building on this insight, we implement a lightweight alignment framework via targeted activation steering within these ToM-critical layers.Experiments demonstrate that COSTOM significantly enhances human-like social reasoning capabilities and downstream dialogue quality. Mengfan Li 0001, Xuanhua Shi, Yang Deng 0002 |
ACL (1) | 2 |
| 2026 | HieraNTT: A Memory Hierarchy-Aware Data Access Architecture for Efficient Number Theoretic Transform on GPU
Qian Xiong, Weiliang Ma, Ligang He, Yufan Bai, Yao Chen 0008, Hai Jin 0001, Xuanhua Shi |
APPT | 7 |
| 2026 | ParetoES: Hardware-Accelerated Sparse Embedding Similarity via Pareto-Optimal Pruning
Jiaqi Zhai, Xuanhua Shi, Wenju Zhao, Chencheng Ye 0001, Shunsen Lv, Zhongtian Long, Bingsheng He, Hai Jin 0001 |
ISCA | 2 |
| 2026 | Multi-instance multi-label position-aware doubly graph convolutional networks
Zhi Li 0048, Teng Zhang 0001, Caiwu Jiang, Xuanhua Shi, Hai Jin 0001 |
Frontiers Comput. Sci. | 5 |
| 2025 | AccelES: Accelerating Top-K SpMV for Embedding Similarity via Low-bit PruningabstractIn the realm of recommendation systems, achieving real-time performance in embedding similarity tasks is often hindered by the limitations of traditional Top-K sparse matrix-vector multiplication (SpMV) methods, which suffer from high latency due to inefficient memory access patterns. This paper identifies these critical gaps and introduces AccelES, a novel approach that significantly enhances the efficiency of Top-K SpMV. Our method employs a two-stage calculation scheme: the first stage utilizes a compact, low-bit dataset to quickly identify the most relevant entries, while the second stage performs full-precision calculations solely on this pruned subset, thereby minimizing computational overhead. Furthermore, AccelES incorporates innovative matrix representations, Ultra-CSR and Random-CSR, which optimize memory bandwidth utilization. Experimental results demonstrate that AccelES accelerates performance, surpassing state-of-the-art FPGA, GPU, and CPU solutions by factors of 3.4×, 2.5×, and 153.3×, respectively, under controlled conditions. These advancements not only enhance processing speed but also significantly improve real-time performance in recommendation systems, establishing AccelES as a pivotal contribution to the field of Top-K sparse matrix-vector multiplication. Jiaqi Zhai, Xuanhua Shi, Chencheng Ye 0001, Weifang Hu, Bingsheng He, Hai Jin 0001 |
HPCA | 2 |
| 2025 | Dataflow-Guided Neuro-Symbolic Language Models for Type InferenceabstractLanguage Models (LMs) are increasingly used for type inference, aiding in error detection and software development.
Some real-world deployments of LMs require the model to run on local machines to safeguard the intellectual property of the source code. This setting often limits the size of the LMs that can be used. We present Nester, the first neuro-symbolic approach that enhances LMs for type inference by integrating symbolic learning without increasing model size. Nester breaks type inference into sub-tasks based on the data and control flow of the input code, encoding them as a modular high-level program. This program executes multi-step actions, such as evaluating expressions and analyzing conditional branches of the target code, combining static typing with LMs to infer potential types.
Evaluated on the ManyTypes4Py dataset in Python, Nester outperforms two state-of-the-art type inference methods (HiTyper and TypeGen), achieving 70.7\% Top-1 Exact Match, which is 18.3\% and 3.6\% higher than HiTyper and TypeGen, respectively. For complex type annotations like typing.Optional and typing.Union, Nester achieves 51.0\% and 16.7\%, surpassing TypeGen by 28.3\% and 5.8\%. Ge Li 0001, Yao Wan 0001, Hongyu Zhang 0002, Zhou Zhao 0001, Wenbin Jiang 0001, Xuanhua Shi, Hai Jin 0001, Zheng Wang 0001 |
ICML | 6 |
| 2025 | CodeSync: Synchronizing Large Language Models with Dynamic Code Evolution at ScaleabstractLarge Language Models (LLMs) have exhibited exceptional performance in software engineering yet face challenges in adapting to continually evolving code knowledge, particularly the frequent updates of third-party library APIs. This limitation, rooted in the static pre-training datasets, often results in non-executable code or implementations with suboptimal safety and efficiency. To this end, we introduce CodeSync, a data engine to identify outdated code patterns and collect real-time code knowledge updates from Python third-party libraries. Building upon CodeSync, we develop CodeSyncBench, a comprehensive benchmark for assessing LLMs' ability to stay synchronized with code evolution, which covers real-world updates for 220 APIs from six Python libraries. Our benchmark offers 3,300 test cases spanning three evaluation tasks and an update-aware instruction tuning dataset of 2,200 training samples. Extensive experiments on 14 LLMs reveal that they struggle with dynamic code evolution, even with the support of advanced knowledge updating methods (e.g., DPO, ORPO, and SimPO). Our CodeSync lays a strong foundation for developing more effective and robust methods for real-time code knowledge updating in the future. The experimental code is available at: https://github.com/CGCL-codes/naturalcc/tree/main/examples/codesync. Zhaoyang Chu, Zhengxiang Cheng, Xuyi Yang, Kaiyue Qiu, Yao Wan 0001, Zhou Zhao 0001, Xuanhua Shi, Hai Jin 0001, Dongping Chen |
ICML | 8 |
| 2025 | LaTCoder: Converting Webpage Design to Code with Layout-as-ThoughtabstractConverting webpage designs into code (design-to-code) plays a vital role in User Interface (UI) development for front-end developers, bridging the gap between visual design and functional implementation. While recent Multimodal Large Language Models (MLLMs) have shown significant potential in design-to-code tasks, they often fail to accurately preserve the layout during code generation. To this end, we draw inspiration from the Chain-of-Thought (CoT) reasoning in human cognition and propose LaTCoder, a novel approach that enhances layout preservation in webpage design during code generation with Layout-as-Thought (LaT). Specifically, we first introduce a simple yet efficient algorithm to divide the webpage design into image blocks. Next, we prompt MLLMs using a CoT-based approach to generate code for each block. Finally, we apply two assembly strategies-absolute positioning and an MLLM-based method-followed by dynamic selection to determine the optimal output. We evaluate the effectiveness of LaTCoder using multiple backbone MLLMs (i.e., DeepSeek-VL2, Gemini, and GPT-4o) on both a public benchmark and a newly introduced, more challenging benchmark (CC-HARD) that features complex layouts. The experimental results on automatic metrics demonstrate significant improvements. Specifically, TreeBLEU scores increased by 66.67% and MAE decreased by 38% when using DeepSeek-VL2, compared to direct prompting. Moreover, the human preference evaluation results indicate that annotators favor the webpages generated by LaTCoder in over 60% of cases, providing strong evidence of the effectiveness of our method. Yi Gui, Zhen Li 0050, Guohao Wang, Tianpeng Lv, Gaoyang Jiang, Yi Liu 0069, Dongping Chen, Yao Wan 0001, Hongyu Zhang 0002, Wenbin Jiang 0001, Xuanhua Shi, Hai Jin 0001 |
KDD (2) | 12 |
| 2025 | RE-SEGNN: recurrent semantic evidence-aware graph neural network for temporal knowledge graph forecasting
Wenyu Cai, Mengfan Li 0001, Xuanhua Shi, Yuanxin Fan, Quntao Zhu, Hai Jin 0001 |
Sci. China Inf. Sci. | 3 |
| 2025 | E2CNN: entity-type-enriched cascaded neural network for Chinese financial relation extractionabstractAbstract Knowledge Graphs (KGs) are pivotal for effectively organizing and managing structured information across various applications. Financial KGs have been successfully employed in advancing applications such as audit, anti-fraud, and anti-money laundering. Despite their success, the construction of Chinese financial KGs has seen limited research due to the complex semantics. A significant challenge is the overlap triples problem, where entities feature in multiple relations within a sentence, hampering extraction accuracy–more than 39% of the triples in Chinese datasets exhibit the overlap triples. To address this, we propose the Entity-type-Enriched Cascaded Neural Network (E 2 CNN), leveraging special tokens for entity boundaries and types. E 2 CNN ensures consistency in entity types and excludes specific relations, mitigating overlap triple problems and enhancing relation extraction. Besides, we introduce the available Chinese financial dataset F in C orpus .CN, annotated from annual reports of 2,000 companies, containing 48,389 entities and 23,368 triples. Experimental results on the DUIE dataset and F in C orpus .CN underscore E 2 CNN’s superiority over state-of-the-art models. Mengfan Li 0001, Xuanhua Shi, Chenqi Qiao, Yao Wan 0001, Teng Zhang 0001, Hai Jin 0001 |
Frontiers Comput. Sci. | 2 |
| 2025 | Text-augmented long-term relation dependency learning for knowledge graph representationabstractKnowledge graph (KG) representation learning aims to map entities and relations into a low-dimensional representation space, showing significant potential in many tasks. Existing approaches follow two categories: (1) Graph-based approaches encode KG elements into vectors using structural score functions. (2) Text-based approaches embed text descriptions of entities and relations via pre-trained language models (PLMs), further fine-tuned with triples. We argue that graph-based approaches struggle with sparse data, while text-based approaches face challenges with complex relations. To address these limitations, we propose a unified Text-Augmented Attention-based Recurrent Network, bridging the gap between graph and natural language. Specifically, we employ a graph attention network based on local influence weights to model local structural information and utilize a PLM based prompt learning to learn textual information, enhanced by a mask-reconstruction strategy based on global influence weights and textual contrastive learning for improved robustness and generalizability. Besides, to effectively model multi-hop relations, we propose a novel semantic-depth guided path extraction algorithm and integrate cross-attention layers into recurrent neural networks to facilitate learning the long-term relation dependency and offer an adaptive attention mechanism for varied-length information. Extensive experiments demonstrate that our model exhibits superiority over existing models across KG completion and question-answering tasks. Quntao Zhu, Mengfan Li 0001, Yuanjun Gao, Yao Wan 0001, Xuanhua Shi, Hai Jin 0001 |
High Confid. Comput. | 5 |
| 2025 | gECC: A GPU-based high-throughput framework for Elliptic Curve CryptographyabstractElliptic Curve Cryptography (ECC) is an encryption method that provides security comparable to traditional techniques like Rivest–Shamir–Adleman (RSA) but with lower computational complexity and smaller key sizes, making it a competitive option for applications such as blockchain, secure multi-party computation, and database security. However, the throughput of ECC is still hindered by the significant performance overhead associated with elliptic curve (EC) operations, which can affect their efficiency in real-world scenarios. This article presents gECC , a versatile framework for ECC optimized for GPU architectures, specifically engineered to achieve high-throughput performance in EC operations. To maximize throughput, gECC incorporates batch-based execution of EC operations and microarchitecture-level optimization of modular arithmetic. It employs Montgomery’s trick [ 40 ] to enable batch EC computation and incorporates novel computation parallelization and memory management techniques to maximize the computation parallelism and minimize the access overhead of GPU global memory. Furthermore, we analyze the primary bottleneck in modular multiplication by investigating how the user codes of modular multiplication are compiled into hardware instructions and what these instructions’ issuance rates are. We identify that the efficiency of modular multiplication is highly dependent on the number of Integer Multiply-Add (IMAD) instructions. To eliminate this bottleneck, we propose novel techniques to minimize the number of IMAD instructions by leveraging predicate registers to pass the carry information and using addition and subtraction instructions (IADD3) to replace IMAD instructions. Our experimental results show that, for ECDSA and ECDH, the two commonly used ECC algorithms, gECC can achieve performance improvements of 5.56 × and 4.94 ×, respectively, compared to the state-of-the-art GPU-based system. In a real-world blockchain application, we can achieve performance improvements of 1.56 ×, compared to the state-of-the-art CPU-based system. gECC is completely and freely available at https://github.com/CGCL-codes/gECC . Qian Xiong, Weiliang Ma, Xuanhua Shi, Yongluan Zhou, Hai Jin 0001, Haozhou Wang, Zhengru Wang |
ACM Trans. Archit. Code Optim. | 3 |
| 2025 | Corrigendum: gECC: A GPU-based high-throughput framework for Elliptic Curve CryptographyabstractThis is a corrigendum for the article “gECC: A GPU-based high-throughput framework for Elliptic Curve Cryptography” published in ACM Trans. Arch. Code Optim. 22, 3, Article 84 (September 2025), 27 pages. Qian Xiong, Weiliang Ma, Xuanhua Shi, Yongluan Zhou, Hai Jin 0001, Haozhou Wang, Zhengru Wang |
ACM Trans. Archit. Code Optim. | 3 |
| 2025 | RuYi: Optimizing Burst Buffer Through Automated, Fine-Grained Process-to-BB MappingabstractCurrent supercomputers use an SSD-based storage layer called Burst Buffer (BB) to provide I/O-intensive applications with accelerated storage access. However, efficiently utilizing this limited and expensive storage remains a critical issue, creating an urgent need for implementing Quality of Service (QoS) in BB. To address this, we propose RuYi, a QoS-aware method to provide applications with bandwidth guarantees in the BB file system. RuYi tackles two main issues. First, it quantitatively profiles available bandwidth resources in BB to ensure reliable QoS, a crucial aspect seldom studied in the literature. Second, RuYi offers fine-grained process-level QoS via an innovative process-to-BB mapping, maximizing resource utilization—something not achievable with conventional coarse-grained compute-to-BB mapping. We evaluated RuYi on a subsystem of the leading exascale supercomputer Sunway, consisting of 4,000 compute nodes and 200 BB nodes. The experimental results demonstrate that RuYi achieves an impressive end-to-end bandwidth control accuracy of 97%, while improving BB utilization by up to 116% compared to conventional coarse-grained compute-to-BB mapping. Yusheng Hua, Xuanhua Shi, Ligang He, Teng Zhang 0001, Hai Jin 0001, Yong Chen 0001 |
IEEE Trans. Computers | 2 |
| 2025 | Can Large Language Models Serve as Evaluators for Code Summarization?abstractCode summarization facilitates program comprehension and software maintenance by converting code snippets into natural-language descriptions. Over the years, numerous methods have been developed for this task, but a key challenge remains: effectively evaluating the quality of generated summaries. While human evaluation is effective for assessing code summary quality, it is labor-intensive and difficult to scale. Commonly used automatic metrics, such as BLEU, ROUGE-L, METEOR, and BERTScore, often fail to align closely with human judgments. In this paper, we explore the potential ofLarge Language Models (LLMs)for evaluating code summarization. We propose CODERPE (Role-Player for Code Summarization Evaluation), a novel method that leverages role-player prompting to assess the quality of generated summaries. Specifically, we prompt LLM-based evaluators to take on diverse roles, such as code reviewer, code author, code editor, and system analyst. Each role evaluates the quality of code summaries across key dimensions, including coherence, consistency, fluency, and relevance. We further explore the robustness of LLMs as evaluators by employing various prompting strategies, including chain-of-thought reasoning, incontext learning, and tailored rating form designs. The results demonstrate that LLMs serve as effective evaluators for code summarization. Notably, our LLM-based evaluator, CODERPE , achieves an 80.18% Spearman correlation with human evaluations, outperforming the existing BERTScore metric by 10.39%. Yang Wu 0010, Yao Wan 0001, Zhaoyang Chu, Wenting Zhao 0006, Ye Liu 0006, Hongyu Zhang 0002, Xuanhua Shi, Hai Jin 0001, Philip S. Yu |
IEEE Trans. Software Eng. | 7 |
| 2024 | Towards Understanding the Effectiveness of Large Language Models on Directed Test Input GenerationabstractAutomatic testing has garnered significant attention and success over the past few decades. Techniques such as unit testing and coverage-guided fuzzing have revealed numerous critical software bugs and vulnerabilities. However, a long-standing, formidable challenge for existing techniques is how to achieve higher testing coverage. Constraint-based techniques, such as symbolic execution and concolic testing, have been well-explored and integrated into the existing approaches. With the popularity of Large Language Models (LLMs), recent research efforts to design tailored prompts to generate inputs that can reach more uncovered target branches. However, the effectiveness of using LLMs for generating such directed inputs and the comparison with the proven constraint-based solutions has not been systematically explored. Zongze Jiang, Ming Wen 0001, Jialun Cao, Xuanhua Shi, Hai Jin 0001 |
ASE | 4 |
| 2024 | Uncovering Nested Data Parallelism and Data Reuse in DNN Computation with FractalTensorabstractTo speed up computation, deep neural networks (DNNs) usually rely on highly optimized tensor operators. Despite the effectiveness, tensor operators are often defined empirically with ad hoc semantics. This hinders the analysis and optimization across operator boundaries. FractalTensor is a programming framework that addresses this challenge. At the core, FractalTensor is a nested list-based abstract data type (ADT), where each element is a tensor with static shape or another FractalTensor (i.e., nested). DNNs are then de-fined by high-order array compute operators like map/reduce/scan and array access operators like window/stride on FractalTensor. This new way of DNN definition explicitly exposes nested data parallelism and fine-grained data access patterns, opening new opportunities for whole program analysis and optimization. To exploit these opportunities, from the FractalTensor-based code the compiler extracts a nested multi-dimensional dataflow graph called Extended Task Dependence Graph (ETDG), which provides a holistic view of data dependency across different granularity. The ETDG is then transformed into an efficient implementation through graph coarsening, data reordering, and access materialization. Evaluation on six representative DNNs like RNN and FlashAttention on NVIDIA A100 shows that Fractal-Tensor achieves speedup by up to 5.45x and 2.14x on average through a unified solution for diverse optimizations. Siran Liu, Chengxiang Qi, Chao Yang 0002, Weifang Hu, Xuanhua Shi, Fan Yang 0024, Mao Yang 0004 |
SOSP | 6 |
| 2024 | QoS-pro: A QoS-enhanced Transaction Processing Framework for Shared SSDsabstractSolid State Drives (SSDs) are widely used in data-intensive scenarios due to their high performance and decreasing cost. However, in shared environments, concurrent workloads can interfere with each other, leading to a violation of Quality of Service (QoS). While QoS mechanisms like fairness guarantees and latency constraints have been integrated into SSDs, existing transaction processing frameworks offer limited QoS guarantees and can significantly degrade overall performance in a shared environment. The reason is that the internal components of an SSD, originally designed to exploit parallelism, struggle to coordinate effectively when QoS mechanisms are applied to them. This article proposes a novel QoS -enhanced transaction pro cessing framework, called QoS-pro, which enhances QoS guarantees for concurrent workloads while maintaining high parallelism for SSDs. QoS-pro achieves this by redesigning transaction processing procedures to fully exploit the parallelism of shared SSDs and enhancing QoS-oriented transaction translation and scheduling with parallelism features in mind. In terms of fairness guarantees, QoS-pro outperforms state-of-the-art methods by achieving 96% fairness improvement and 64% maximum latency reduction. QoS-pro also shows almost no loss in throughput when compared with parallelism-oriented methods. Additionally, QoS-pro triggers the fewest Garbage Collection (GC) operations and minimally affects concurrently running workloads during GC operations. Hao Fan 0006, Yiliang Ye, Shadi Ibrahim, Xingru Li, Weibin Xue, Song Wu 0001, Chen Yu 0003, Xuanhua Shi, Hai Jin 0001 |
ACM Trans. Archit. Code Optim. | 9 |
| 2024 | Core Maintenance on Dynamic Graphs: A Distributed Approach Built on H-IndexabstractCore number is an essential tool for analyzing graph structure. Graphs in the real world are typically large and dynamic, requiring the development of distributed algorithms to refrain from expensive I/O operations and the maintenance algorithms to address dynamism. Core maintenance updates the core number of each vertex upon the insertion/deletion of vertices/edges. Although the state-of-the-art distributed maintenance algorithm [9] can handle multiple edge insertions/deletions simultaneously, it still has two aspects to improve. (I) Parallel processing is not allowed when inserting/removing edges with the same core number, reducing the degree of parallelism and raising the number of rounds. (II) During the implementation phase, only one thread is assigned to the vertices with the same core number, leading to the inability to fully utilize the distributed computing power. Furthermore, the h-index [1] based distributed core decomposition algorithm [10] can fully utilize the distributed computing power where all vertices can be processed in parallel. However, it requires all vertices to recompute their core numbers upon graph changes. In this article, we propose a distributed core maintenance algorithm based on h-index, which circumvents the issues of algorithm [9]. In addition, our algorithm avoids core numbers recalculation where the numbers do not change. In comparison to the state-of-the-art distributed maintenance algorithm [9], the time speedup ratio is at least 100 in the scenarios of both insertion and deletion. Compared to the distributed core decomposition algorithm [10], the average time speedup ratios are 2 and 8 for the cases of insertion and deletion, respectively. Qiang-Sheng Hua, Hongen Wang, Hai Jin 0001, Xuanhua Shi |
IEEE Trans. Big Data | 4 |
| 2024 | MMDataLoader: Reusing Preprocessed Data Among Concurrent Model Training TasksabstractData preprocessing plays an important role in deep learning, which directly affects the training efficiency. Data preprocessing is performed on the CPU. The preprocessed data are then fed to the models that are trained on the GPU. We observe that data preprocessing on the CPU can potentially create a bottleneck in the entire process of a model training task. In order to tackle this issue, we have developed MMDataLoader, which enables reusing preprocessed data among multiple model training tasks. MMDataLoader automatically constructs a data preprocessing pipeline based on each task's specific preprocessing workflow, allowing for maximum data reuse and reduced computing workload on the CPU. Unlike conventional data loaders that operate at the task level and provide data provision services to specific training tasks, MMDataLoader operates at the server level and provides data for all concurrently running tasks. We have conducted extensive experiments. The results show that MMDataLoader can significantly increase preprocessing throughput without affecting model convergence when compared to conventional methods where model training tasks are executed concurrently. For instance, with three tasks running, the preprocessing throughput can increase by 1.6x to 3.15x, depending on the tasks being executed and the proportion of preprocessing operations that are shared among them. Hai Jin 0001, Zhanyang Zhu, Ligang He, Yusheng Hua, Xuanhua Shi |
IEEE Trans. Computers | 6 |
| 2023 | GZKP: A GPU Accelerated Zero-Knowledge Proof SystemabstractZero-knowledge proof (ZKP) is a cryptographic protocol that allows one party to prove the correctness of a statement to another party without revealing any information beyond the correctness of the statement itself. It guarantees computation integrity and confidentiality, and is therefore increasingly adopted in industry for a variety of privacy-preserving applications, such as verifiable outsource computing and digital currency. Weiliang Ma, Qian Xiong, Xuanhua Shi, Xiaosong Ma, Hai Jin 0001, Haozhao Kuang, Mingyu Gao 0001, Ye Zhang 0042, Haichen Shen, Weifang Hu |
ASPLOS (2) | 3 |
| 2023 | Incremental and Decremental Optimal Margin Distribution LearningabstractIncremental and decremental learning (IDL) deals with the tasks where new data arrives sequentially as a stream or old data turns unavailable continually due to the privacy protection. Existing IDL methods mainly focus on support vector machine and its variants with linear-type loss. There are few studies about the quadratic-type loss, whose Lagrange multipliers are unbounded and much more difficult to track. In this paper, we take the latest statistical learning framework optimal margin distribution machine (ODM) which involves a quadratic-type loss due to the optimization of margin variance, for example, and equip it with the ability to handle IDL tasks. Our proposed ID-ODM can avoid updating the Lagrange multipliers in an infinite range by determining their optimal values beforehand so as to enjoy much more efficiency. Moreover, ID-ODM is also applicable when multiple instances come and leave simultaneously. Extensive empirical studies show that ID-ODM can achieve 9.1x speedup on average with almost no generalization lost compared to retraining ODM on new data set from scratch. Li-Jun Chen, Teng Zhang 0001, Xuanhua Shi, Hai Jin 0001 |
IJCAI | 3 |
| 2023 | Scalable Optimal Margin Distribution MachineabstractOptimal margin Distribution Machine (ODM) is a newly proposed statistical learning framework rooting in the novel margin theory, which demonstrates better generalization performance than the traditional large margin based counterparts. Nonetheless, it suffers from the ubiquitous scalability problem regarding both computation time and memory as other kernel methods. This paper proposes a scalable ODM, which can achieve nearly ten times speedup compared to the original ODM training method. For nonlinear kernels, we propose a novel distribution-aware partition method to make the local ODM trained on each partition be close and converge faster to the global one. When linear kernel is applied, we extend a communication efficient SVRG method to accelerate the training further. Extensive empirical studies validate that our proposed method is highly computational efficient and almost never worsen the generalization. Nan Cao 0002, Teng Zhang 0001, Xuanhua Shi, Hai Jin 0001 |
IJCAI | 4 |
| 2023 | HyperBit: A temporal graph store for fast answering queries
Shaoqi Zang, Pingpeng Yuan, Xuanhua Shi, Hai Jin 0001 |
Data Knowl. Eng. | 4 |
| 2023 | Data Stream Clustering: An In-depth Empirical StudyabstractData Stream Clustering (DSC) plays an important role in mining continuous and unlabeled data streams in real-world applications. Over the last decades, numerous DSC algorithms have been proposed with promising clustering accuracy and efficiency. Despite the significant differences among existing DSC algorithms, they are commonly built around four key design aspects: summarizing data structure, window model, outlier detection mechanism, and offline refinement strategy. However, there is a lack of empirical studies on these key design aspects in the same codebase using real-world workloads with distinct characteristics. As a result, it is difficult for researchers to improve upon the state-of-the-art. In this paper, we conduct such a study of DSC on its four key design aspects. We implemented state-of-the-art variants of all of these design choices in an open-sourced platform from scratch and evaluated them using both real-world and synthetic workloads. Our analysis identifies the fundamental issues and trade-offs of each design choice in terms of both accuracy and efficiency. We even find that combining flexible design choices led to the development of a new algorithm called Benne, which can be tuned to achieve either better accuracy or better efficiency compared to the state-of-the-art. Xin Wang 0120, Zhengru Wang, Shuhao Zhang 0001, Xuanhua Shi |
Proc. ACM Manag. Data | 5 |
| 2023 | Parallel Overlapping Community Detection Algorithm on GPUabstractCommunity detection is one of the most representative graph mining applications, which is often assembled as a concurrent graph partition application to explore the maximum modularity (or gained modularity) of each community. However, many branch divergence operations create significant obstacles to unleashing GPU's high throughput and memory bandwidth, which are needed in community detection applications to divide the vertices into different communities. In this paper, we present Lugger, a GPU-based overlapping community detection algorithm that reduces GPU's branch divergence via the customer-designed cache-aware parallel searching technique. In Lugger, we first design a cache-aware parallel searching policy using the B-Tree structure. Then, we set the B-Tree node matches with the GPU cache line to meet the coalesced memory access manner and avoid the branch divergence in warps. Moreover, we design a positive node splitting scheme to reduce the lock operation and idle threads when building the B-Tree structure. In addition, we implement a warp-centric thread assignment strategy to make sure the workloads across threads are balanced. We implement the proposed algorithm on NVIDIA GPU and evaluate the performance on eight large graphs (up to$\text{3}~M$vertices and$\text{117}~M$edges) with ground-truth communities. The experimental results show that Lugger can outperform the state-of-the-art works on scalability and detection quality. Zhigao Zheng 0001, Xuanhua Shi, Hai Jin 0001 |
IEEE Trans. Big Data | 2 |
| 2023 | Waterwave: A GPU Memory Flow Engine for Concurrent DNN TrainingabstractTraining Deep Neural Networks (DNN) concurrently is becoming increasingly important for deep learning practitioners, e.g.,hyperparameter optimization (HPO)andneural architecture search (NAS). The GPU memory capacity is the impediment that prohibits multiple DNNs from being trained on the same GPU due to the large memory usage during training. In this paper, we proposeWaterwave, a GPU memory flow engine for concurrent deep learning training. First, to address the memory explosion brought by the long time lag between memory allocation and deallocation time, we develop an allocator tailored for multi-streams. By making the allocator aware of the stream information, aprioritized allocationis conducted based on the chunk'ssynchronizationattributes, allowing us to provide useable memory after scheduling rather than waiting it to be really released after GPU computation. Second,Waterwavepartitions the compute graph to a set of continuousnode groupsand then performs finer-grained scheduling:NodeGroup pipeline execution, to guarantee a proper memory requests order.Waterwavecan accomplish up to 96.8% of the maximum batch size of solo training. Additionally, in scenarios with high memory demand,Waterwavecan outperform existing spatial sharing and temporal sharing by up to 12x and 1.49x, respectively. Xuanhua Shi, Ligang He, Yunfei Zhao 0001, Hai Jin 0001 |
IEEE Trans. Computers | 1 |
| 2023 | TurboGNN: Improving the End-to-End Performance for Sampling-Based GNN Training on GPUsabstractGraph Neural Networks(GNN) have evolved as powerful models for graph representation learning. Sampling-based training methods have been introduced to train large graphs without compromising accuracy. However, it is challenging for the existing GNN systems to effectively utilize multi-core accelerators, especially GPUs, due to a large number of atomic operations and unbalanced workload originating from the serial execution of multiple GNN processing stages. In this paper, we propose a combination of optimization techniques to accelerate the end-to-end performance of the sampling-based GNN training process. Specifically, we propose an adaptive share memory-based sampling technique and a degree-guided thread block scheduling strategy to optimize the graph sampling. Further, based on the observations of resource demand in different training stages, we propose an asynchronous pipeline-based scheduling method, which accelerates the GNN training by decoupling different training stages into a pipeline and therefore improves the GPU resource utilization significantly. The experimental results show that compared with the existing work, the proposed methods can achieve up to 5.6X performance speedup in the end-to-end performance. Xuanhua Shi, Ligang He, Hai Jin 0001 |
IEEE Trans. Computers | 2 |
| 2023 | Temporal Graph CubeabstractData warehouse and OLAP (Online Analytical Processing) are effective tools for decision support on traditional relational data and static multidimensional network data. However, many real-world multidimensional networks are often modeled as temporal multidimensional networks, where the edges in the network are associated with temporal information. Such temporal multidimensional networks typically cannot be handled by traditional data warehouse and OLAP techniques. To fill this gap, we propose a novel data warehouse model, named$\mathsf {Temporal{ }\; Graph{ }\; Cube}$, to support OLAP queries on temporal multidimensional networks. Through supporting OLAP queries in any time range, users can obtain summarized information of the network in the time range of interest, which cannot be derived by using traditional static graph OLAP techniques. We propose a segment-tree based indexing technique to speed up the OLAP queries, and also develop an index-updating technique to maintain the index when the temporal multidimensional network evolves over time. In addition, we also propose a novel concept called$\mathsf {similarity{ }\; of{ }\; snapshots}$which shows a strong correlation with the efficiency of indexing technique and can provide a good reference on the necessity of building the index. The results of extensive experiments on two large real-world datasets demonstrate the effectiveness and efficiency of the proposed method. Guoren Wang, Yue Zeng 0004, Rong-Hua Li 0001, Hongchao Qin, Xuanhua Shi, Yubin Xia, Xuequn Shang 0001, Liang Hong 0001 |
IEEE Trans. Knowl. Data Eng. | 5 |
| 2023 | TurboMGNN: Improving Concurrent GNN Training Tasks on GPU With Fine-Grained Kernel FusionabstractGraph Neural Networks(GNN) have evolved as powerful models for graph representation learning. Many works have been proposed to support GNN training efficiently on GPU. However, these works only focus on a single GNN training task such as operator optimization, task scheduling, and programming model. Concurrent GNN training, which is needed in the applications such as neural network structure search, has not been explored yet. This work aims to improve the training efficiency of the concurrent GNN training tasks on GPU by developing fine-grained methods to fuse the kernels from different tasks. Specifically, we propose a fine-grainedSparse Matrix Multiplication(SpMM) based kernel fusion method to eliminate redundant accesses to graph data. In order to increase the fusion opportunity and reduce the synchronization cost, we further propose a novel technique to enable the fusion of the kernels in forward and backward propagation. Finally, in order to reduce the resource contention caused by the increased number of concurrent, heterogeneous GNN training tasks, we propose an adaptive strategy to group the tasks and match their operators according to resource contention. We have conducted extensive experiments, including kernel- and model-level benchmarks. The results show that the proposed methods can achieve up to 2.6X performance speedup. Xuanhua Shi, Ligang He, Hai Jin 0001 |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2022 | Posistive-Unlabeled Learning via Optimal Transport and Margin DistributionabstractPositive-unlabeled (PU) learning deals with the circumstances where only a portion of positive instances are labeled, while the rest and all negative instances are unlabeled, and due to this confusion, the class prior can not be directly available. Existing PU learning methods usually estimate the class prior by training a nontraditional probabilistic classifier, which is prone to give an overestimation. Moreover, these methods learn the decision boundary by optimizing the minimum margin, which is not suitable in PU learning due to its sensitivity to label noise. In this paper, we enhance PU learning methods from the above two aspects. More specifically, we first explicitly learn a transformation from unlabeled data to positive data by entropy regularized optimal transport to achieve a much more precise estimation for class prior. Then we switch to optimizing the margin distribution, rather than the minimum margin, to obtain a label noise insensitive classifier. Extensive empirical studies on both synthetic and real-world data sets demonstrate the superiority of our proposed method. Nan Cao 0002, Teng Zhang 0001, Xuanhua Shi, Hai Jin 0001 |
IJCAI | 3 |
| 2022 | Enhance Temporal Knowledge Graph Completion via Time-Aware Attention Graph Convolutional Network
HaoHui Wei, Hong Huang 0001, Teng Zhang 0001, Xuanhua Shi, Hai Jin 0001 |
ECML/PKDD (2) | 4 |
| 2022 | Reveal training performance mystery between TensorFlow and PyTorch in the single GPU environment
Hulin Dai, Xuanhua Shi, Ligang He, Qian Xiong, Hai Jin 0001 |
Sci. China Inf. Sci. | 3 |
| 2022 | LoomIO: Object-Level Coordination in Distributed File SystemsabstractDevice-level interference is recognized as a major cause of the performance degradation in distributed file systems. Although the approaches of mitigating interference through coordination at application-level, middleware-level, and server-level have shown beneficial results in previous studies, we find their effectiveness is largely reduced since I/O requests are re-arranged by underlying object file systems. In this research study, we prove that object-level coordination is critical and often the key to address the interference issue, as the scheduling of object requests determines the device-level accesses and thus determines the actual I/O bandwidth and latency. This article proposes an object-level coordination system, LoomIO, which uses an OBOP (One-Broadcast-One-Propagate) method and a time-limited coordination process to deliver highly efficient coordination service. Specifically, LoomIO enables object requests to achieve an optimized scheduling decision within a few milliseconds and largely mitigates the device-level interference. We have implemented a LoomIO prototye and integrated it into Ceph file system. The evaluation results show that LoomIO achieved the considerable improvements in resource utilization (by up to 35%), in I/O throughput (by up to 31%), and in 99th percentile latency (by up to 54%) compared to the K-optimal method which uses the same scheduling algorithm as LoomIO but does not have the coordination support. Yusheng Hua, Xuanhua Shi, Hai Jin 0001, Wei Xie 0017, Ligang He, Yong Chen 0001 |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2022 | Modeling Information Diffusion With Sequential Interactive HypergraphsabstractIn online social networks, numerous users generate and distribute tremendous content simultaneously, and how information spreads in online society has attracted critical attention. Recent studies introduce deep learning to help mine the complex diffusion pattern and forecast the diffusion trend. However, there are two limitations of the state-of-art methods. On the one hand, existing models only focus on the internal influence within the diffusion flow but ignore the external influence from the dissemination of other contents. On the other hand, the dynamics of user interest are barely considered, while the changes of user preference for contents also have a considerable impact on future diffusion. To address these issues, we introduce the hypergraph structure and a sequential framework to model complex interactions in social networks. Then, we propose dual-channel hypergraph neural networks to tackle the diffusion prediction problem, denoted by HyperINF. Specifically, in the user channel, we build sequential user interactive hypergraphs and learn the dynamic user representation, and in the diffusion channel, we construct a diffusion interactive graph to capture the cross-diffusion relation. At last, we consider the social relation to help make the prediction. Experimental results on three datasets suggest the effectiveness and practicability of the proposed framework. Hai Jin 0001, Hong Huang 0001, Yu Song 0005, HaoHui Wei, Xuanhua Shi |
IEEE Trans. Sustain. Comput. | 6 |
| 2021 | A nearly optimal distributed algorithm for computing the weighted girth
Qiang-Sheng Hua, Lixiang Qian, Dongxiao Yu, Xuanhua Shi, Hai Jin 0001 |
Sci. China Inf. Sci. | 4 |
| 2021 | DDL-QoS: A dynamic I/O scheduling strategy of QoS for HPC applicationsabstractSummary With the increasing cloud‐trend of high‐performance computing (HPC), more users submit their applications simultaneously to the platform and wish they could finish before the deadline. Moreover, due to the severe holistic performance degradation caused by I/O contention, a deadline‐sensitive I/O scheduler is needed to allocate storage resources according to the requirements of applications and resultantly guarantee the quality of service (QoS) of concurrently running applications. In this paper, we first explore the bandwidth allocation phenomenon caused by interference in applications through the modeling of historical data, and then we quote a metric called random percentage that can represent the random degree of the applications and be used to guide I/O scheduling in the later stage. We design a dynamic I/O scheduler named DDL‐QoS that uses solid state drives(SSDs) as QoS guarantee to minimize interference and ensure applications meet their deadline. The potential of our design is that the greater the I/O interference, the greater the performance improvement, but this performance improvement will be limited by the physical properties of the storage hardware. Xuanhua Shi, Wei Liu 0004, Hai Jin 0001, Yusheng Hua |
Concurr. Comput. Pract. Exp. | 2 |
| 2021 | China in the eyes of news media: a case study under COVID-19 epidemicabstractAs one of the early COVID-19 epidemic outbreak areas, China attracted the global news media’s attention at the beginning of 2020. During the epidemic period, Chinese people united and actively fought against the epidemic. However, in the eyes of the international public, the situation reported about China is not optimistic. To better understand how the international public portrays China, especially during the epidemic, we present a case study with big data technology. We aim to answer three questions: (1) What has the international media focused on during the COVID-19 epidemic period? (2) What is the media’s tone when they report China? (3) What is the media’s attitude when talking about China? In detail, we crawled more than 280 000 pieces of news from 57 mainstream media agencies in 22 countries and made some interesting observations. For example, international media paid more attention to Chinese livelihood during the COVID-19 epidemic period. In March and April, “progress of Chinese vaccines,” “specific drugs and treatments,” and “virus outbreak in U.S.” became the media’s most common topics. In terms of news attitude, Cuba, Malaysia, and Venezuela had a positive attitude toward China, while France, Canada, and the United Kingdom had a negative attitude. Our study can help understand China’s image in the eyes of the international media and provide a sound basis for image analysis. Hong Huang 0001, Zhexue Chen, Xuanhua Shi, Zepeng He, Hai Jin 0001, Zongya Li |
Frontiers Inf. Technol. Electron. Eng. | 3 |
| 2021 | TurboDL: Improving the CNN Training on GPU With Fine-Grained Multi-Streaming SchedulingabstractGraphics Processing Units (GPUs) have evolved as powerful co-processors for the CNN training. Many new features have been introduced into GPUs such as concurrent kernel execution and hyper-Q technology. It is challenging to orchestrate concurrency for CNN (convolutional neural networks) training on GPUs since it may introduce synchronization overhead and poor resource utilization. Unlike previous research which mainly focuses on single layer or coarse-grained optimization, we introduce a critical-path based, asynchronous parallelization mechanism, and propose the optimization technique for the CNN training that takes into account global network architecture and GPU resource usage together. The proposed methods can effectively overlap the synchronization and the computation in different streams. As a result, the training process of CNN is accelerated. We have integrated our methods into Caffe. The experimental results show that the Caffe integrated with our methods can achieve 1.30X performance speedup on average compared with Caffe+cuDNN, and even higher performance speedup can be achieved for deeper, wider, and more complicated networks. Hai Jin 0001, Xuanhua Shi, Ligang He, Bing Bing Zhou |
IEEE Trans. Computers | 3 |
| 2021 | Multi-Stage Network Embedding for Exploring Heterogeneous EdgesabstractThe relationships between objects in a network are typically diverse and complex, leading to the heterogeneous edges with different semantic information. In this article, we focus on exploring the heterogeneous edges for network representation learning. By considering each relationship as a view that depicts a specific type of proximity between nodes, we propose a multi-stage non-negative matrix factorization (MNMF) model, committed to utilizing abundant information in multiple views to learn robust network representations. In fact, most existing network embedding methods are closely related to implicitly factorizing the complex proximity matrix. However, the approximation error is usually quite large, since a single low-rank matrix is insufficient to capture the original information. Through a multi-stage matrix factorization process motivated by gradient boosting, our MNMF model achieves lower approximation error. Meanwhile, the multi-stage structure of MNMF gives the feasibility of designing two kinds of non-negative matrix factorization (NMF) manners to preserve network information better. The united NMF aims to preserve the consensus information between different views, and the independent NMF aims to preserve unique information of each view. Concrete experimental results on realistic datasets indicate that our model outperforms three types of baselines in practical applications. Hong Huang 0001, Yu Song 0005, Fanghua Ye 0001, Xing Xie 0001, Xuanhua Shi, Hai Jin 0001 |
ACM Trans. Knowl. Discov. Data | 5 |
| 2021 | Feluca: A Two-Stage Graph Coloring Algorithm With Color-Centric Paradigm on GPUabstractThere are great challenges in performing graph coloring on GPU in general. First, the long-tail problem exists in the recursion algorithm because the conflict (i.e., different threads assign the adjacent nodes to the same color) becomes more likely to occur as the number of iterations increases. Second, it is hard to parallelize the sequential spread algorithm because the color allocation depends on the adjoining iteration. Third, the atomic operation is widely used on GPU to maintain the color list, which can greatly reduce the efficiency of GPU threads. In this article, we propose a two-stage high-performance graph coloring algorithm, called Feluca, aiming to address the above challenges. Feluca combines the recursion-based method with the sequential spread-based method. In the first stage, Feluca uses a recursive routine to color a majority of vertices in the graph. Then, it switches to the sequential spread method to color the remaining vertices in order to avoid the conflicts of the recursive algorithm. Moreover, the following techniques are proposed to further improve the graph coloring performance. i) A new method is proposed to eliminate the cycles in the graph; ii) a top-down scheme is developed to avoid the atomic operation originally required for color selection; and iii) a novel color-centric coloring paradigm is designed to improve the degree of parallelism for the sequential spread part. All these newly developed techniques, together with further GPU-specific optimizations such as coalesced memory access, comprise an efficient parallel graph coloring solution in Feluca. We have conducted extensive experiments on NVIDIA GPU. The results show that Feluca can achieve 1.19 - 8.39× speedup over the state-of-the-art algorithms. Zhigao Zheng 0001, Xuanhua Shi, Ligang He, Hai Jin 0001, Shuo Wei, Hulin Dai |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2020 | Capuchin: Tensor-based GPU Memory Management for Deep LearningabstractIn recent years, deep learning has gained unprecedented success in various domains, the key of the success is the larger and deeper deep neural networks (DNNs) that achieved very high accuracy. On the other side, since GPU global memory is a scarce resource, large models also pose a significant challenge due to memory requirement in the training process. This restriction limits the DNN architecture exploration flexibility. Xuanhua Shi, Hulin Dai, Hai Jin 0001, Weiliang Ma, Qian Xiong, Fan Yang 0024, Xuehai Qian |
ASPLOS | 2 |
| 2020 | ByteSeries: an in-memory time series database for large-scale monitoring systemsabstractMonitoring large-scale and complex systems often generates high-dimensional and highly dynamic time series data. In such a scenario, massive metadata has to be maintained to support efficient querying, whose large footprint poses great challenges to in-memory databases. In this paper, we present ByteSeries, an in-memory time series database that is designed specifically for large-scale monitoring systems to manage high-dimensional time series. We start with an analysis of the production data and workload at ByteDance's metric monitoring system, which contains over 10 billion time series dimensions. The observation of high overhead of metadata management in high-dimensional time series data calls for a rethink of time series database systems. Byte-Series's memory structure employs the novel Compressed Inverted Index to effectively compress metadata while maintaining high efficiency for multi-dimensional queries. In addition, an algorithm is proposed to effectively convert data into compressed form without sacrificing the data ingestion throughput. We experimentally evaluate ByteSeries by comparing it with ByteDance's original production system, tsdc, as well as two open-source systems, namely Gorilla and Prometheus. We show that ByteSeries significantly improves over ByteDance's original production system by 1) reducing the memory footprint of metadata by 60% and the whole memory consumption by 50%, and 2) speeding up multi-dimensional queries by 1.8x-10.7x. Xuanhua Shi, Zezhao Feng, Kaixi Li, Yongluan Zhou, Hai Jin 0001, Bingsheng He, Zhijun Ling |
SoCC | 1 |
| 2020 | Modeling Heterogeneous Edges to Represent Networks with Graph Auto-Encoder
Lu Wang 0002, Yu Song 0005, Hong Huang 0001, Fanghua Ye 0001, Xuanhua Shi, Hai Jin 0001 |
DASFAA (2) | 5 |
| 2020 | Maxson: Reduce Duplicate Parsing Overhead on Raw DataabstractJSON is a very popular data format in many applications in Web and enterprise. Recently, many data analytical systems support the loading and querying JSON data. However, JSON parsing can be costly, which dominates the execution time of querying JSON data. Many previous studies focus on building efficient parsers to reduce this parsing cost, and little work has been done on how to reduce the occurrences of parsing. In this paper, we start with a study with a real production workload in Alibaba, which consists of over 3 million queries on JSON. Our study reveals significant temporal and spatial correlations among those queries, which result in massive redundant parsing operations among queries. Instead of repetitively parsing the JSON data, we propose to develop a cache system named Maxson for caching the JSON query results (the values evaluated from JSONPath) for reuse. Specifically, we develop effective machine learning-based predictor with combining LSTM (long shortterm memory) and CRF (conditional random field) to determine the JSONPaths to cache given the space budget. We have implemented Maxson on top of SparkSQL. We experimentally evaluate Maxson and show that 1) Maxson is able to eliminate the most of duplicate JSON parsing overhead, 2) Maxson improves end-to-end workload performance by 1.5-6.5×. Xuanhua Shi, Hong Huang 0001, Hai Jin 0001, Huan Shen, Yongluan Zhou, Bingsheng He, Ruibo Li, Keyong Zhou |
ICDE | 1 |
| 2020 | Dynamic cluster strategy for hierarchical rollback-recovery protocols in MPI HPC applicationsabstractSummary Fault tolerance in parallel computing becomes increasingly important with a significant rise in high‐performance computing systems. Coordinated checkpointing and message logging protocols are commonly used fault tolerance mechanisms for message‐passing applications. However, these mechanisms are insufficient because of their severe drawbacks. Hierarchical rollback‐recovery protocols, combining coordinated checkpointing with message logging, are a better solution. However, such protocols may not obtain the appropriate efficiency because the communication pattern in different stages of applications may vary at runtime. In an effort to improve the efficiency of hierarchical rollback‐recovery protocols, we propose a dynamic cluster strategy to adapt to the runtime variation of communication pattern by using a prediction scheme. Finally, the efficiency and scalability of the dynamic cluster strategy are evaluated using 2 static process partition algorithms on the High‐Performance Linpack benchmark. Xiaofei Liao, Long Zheng 0003, Binsheng Zhang, Yu Zhang 0027, Hai Jin 0001, Xuanhua Shi |
Concurr. Comput. Pract. Exp. | 6 |
| 2020 | A parameter-level parallel optimization algorithm for large-scale spatio-temporal data mining
Xuanhua Shi, Ligang He, Dongxiao Yu, Hai Jin 0001, Chen Yu 0003, Hulin Dai, Zezhao Feng |
Distributed Parallel Databases | 2 |
| 2020 | Optimizing the SSD Burst Buffer by Traffic DetectionabstractCurrently, HPC storage systems still use hard disk drive (HDD) as their dominant storage device. Solid state drive (SSD) is widely deployed as the buffer to HDDs. Burst buffer has also been proposed to manage the SSD buffering of bursty write requests. Although burst buffer can improve I/O performance in many cases, we find that it has some limitations such as requiring large SSD capacity and harmonious overlapping between computation phase and data flushing phase. In this article, we propose a scheme, called SSDUP+. 1 SSDUP+ aims to improve the burst buffer by addressing the above limitations. First, to reduce the demand for the SSD capacity, we develop a novel method to detect and quantify the data randomness in the write traffic. Further, an adaptive algorithm is proposed to classify the random writes dynamically. By doing so, much less SSD capacity is required to achieve the similar performance as other burst buffer schemes. Next, to overcome the difficulty of perfectly overlapping the computation phase and the flushing phase, we propose a pipeline mechanism for the SSD buffer, in which data buffering and flushing are performed in pipeline. In addition, to improve the I/O throughput, we adopt a traffic-aware flushing strategy to reduce the I/O interference in HDD. Finally, to further improve the performance of buffering random writes in SSD, SSDUP+ transforms the random writes to sequential writes in SSD by storing the data with a log structure. Further, SSDUP+ uses the AVL tree structure to store the sequence information of the data. We have implemented a prototype of SSDUP+ based on OrangeFS and conducted extensive experiments. The experimental results show that our proposed SSDUP+ can save an average of 50% SSD space while delivering almost the same performance as other common burst buffer schemes. In addition, SSDUP+ can save about 20% SSD space compared with the previous version of this work, SSDUP, while achieving 20–30% higher I/O throughput than SSDUP. Xuanhua Shi, Wei Liu 0004, Ligang He, Hai Jin 0001, Yong Chen 0001 |
ACM Trans. Archit. Code Optim. | 1 |
| 2019 | Software-defined QoS for I/O in exascale computing
Yusheng Hua, Xuanhua Shi, Hai Jin 0001, Wei Liu 0004, Yong Chen 0001, Ligang He |
CCF Trans. High Perform. Comput. | 2 |
| 2019 | Vectorizing disks blocks for efficient storage system via deep learning
Dong Dai 0001, Forrest Sheng Bao, Xuanhua Shi, Yong Chen 0001 |
Parallel Comput. | 4 |
| 2019 | Parallel computation of hierarchical closeness centrality and applications
Hai Jin 0001, Dongxiao Yu, Qiang-Sheng Hua, Xuanhua Shi |
World Wide Web | 5 |
| 2018 | GRAM: A GPU-Based Property Graph Traversal and Query for HPC Rich Metadata Management
Wenke Li, Xuanhua Shi, Hong Huang 0001, Hai Jin 0001, Dong Dai 0001, Yong Chen 0001 |
NPC | 2 |
| 2018 | Layrub: layer-centric GPU memory reuse and data migration in extreme-scale deep learning systemsabstractGrowing accuracy and robustness of Deep Neural Networks (DNN) models are accompanied by growing model capacity (going deeper or wider). However, high memory requirements of those models make it difficult to execute the training process in one GPU. To address it, we first identify the memory usage characteristics for deep and wide convolutional networks, and demonstrate the opportunities of memory reuse on both intra-layer and inter-layer levels. We then present Layrub, a runtime data placement strategy that orchestrates the execution of training process. It achieves layer-centric reuse to reduce memory consumption for extreme-scale deep learning that cannot be run on one single GPU. Bo Liu 0057, Wenbin Jiang 0001, Hai Jin 0001, Xuanhua Shi |
PPoPP | 4 |
| 2018 | Layer-Centric Memory Reuse and Data Migration for Extreme-Scale Deep Learning on Many-Core ArchitecturesabstractDue to the popularity of Deep Neural Network (DNN) models, we have witnessed extreme-scale DNN models with the continued increase of the scale in terms of depth and width. However, the extremely high memory requirements for them make it difficult to run the training processes on single many-core architectures such as a Graphic Processing Unit (GPU), which compels researchers to use model parallelism over multiple GPUs to make it work. However, model parallelism always brings very heavy additional overhead. Therefore, running an extreme-scale model in a single GPU is urgently required. There still exist several challenges to reduce the memory footprint for extreme-scale deep learning. To address this tough problem, we first identify the memory usage characteristics for deep and wide convolutional networks, and demonstrate the opportunities for memory reuse at both the intra-layer and inter-layer levels. We then present Layrub, a runtime data placement strategy that orchestrates the execution of the training process. It achieves layer-centric reuse to reduce memory consumption for extreme-scale deep learning that could not previously be run on a single GPU. Experiments show that, compared to the original Caffe, Layrub can cut down the memory usage rate by an average of 58.2% and by up to 98.9%, at the moderate cost of 24.1% higher training execution time on average. Results also show that Layrub outperforms some popular deep learning systems such as GeePS, vDNN, MXNet, and Tensorflow. More importantly, Layrub can tackle extreme-scale deep learning tasks. For example, it makes an extra-deep ResNet with 1,517 layers that can be trained successfully in one GPU with 12GB memory, while other existing deep learning systems cannot. Hai Jin 0001, Bo Liu 0057, Wenbin Jiang 0001, Xuanhua Shi, Bingsheng He, Shaofeng Zhao |
ACM Trans. Archit. Code Optim. | 5 |
| 2018 | Frog: Asynchronous Graph Processing on GPU with Hybrid Coloring ModelabstractGPUs have been increasingly used to accelerate graph processing for complicated computational problems regarding graph theory. Many parallel graph algorithms adopt the asynchronous computing model to accelerate the iterative convergence. Unfortunately, the consistent asynchronous computing requires locking or atomic operations, leading to significant penalties/overheads when implemented on GPUs. As such, the coloring algorithm is adopted to separate the vertices with potential updating conflicts, guaranteeing the consistency/correctness of the parallel processing. Common coloring algorithms, however, may suffer from low parallelism because of a large number of colors generally required for processing a large-scale graph with billions of vertices. We propose a light-weight asynchronous processing framework called Frog with a preprocessing/hybrid coloring model. The fundamental idea is based on the Pareto principle (or 80-20 rule) about coloring algorithms as we observed through masses of real-world graph coloring cases. We find that a majority of vertices (about 80 percent) are colored with only a few colors, such that they can be read and updated in a very high degree of parallelism without violating the sequential consistency. Accordingly, our solution separates the processing of the vertices based on the distribution of colors. In this work, we mainly answer three questions: (1) how to partition the vertices in a sparse graph with maximized parallelism, (2) how to process large-scale graphs that cannot fit into GPU memory, and (3) how to reduce the overhead of data transfers on PCIe while processing each partition. We conduct experiments on real-world data (Amazon, DBLP, YouTube, RoadNet-CA, WikiTalk, and Twitter) to evaluate our approach and make comparisons with well-known non-preprocessed (such as Totem, Medusa, MapGraph, and Gunrock) and preprocessed (Cusha) approaches, by testing four classical algorithms (BFS, PageRank, SSSP, and CC). On all the tested applications and datasets, Frog is able to significantly outperform existing GPU-based graph processing systems except Gunrock and MapGraph. MapGraph gets better performance than Frog when running BFS on RoadNet-CA. The comparison between Gunrock and Frog is inconclusive. Frog can outperform Gunrock more than 1.04X when running PageRank and SSSP, while the advantage of Frog is not obvious when running BFS and CC on some datasets especially for RoadNet-CA. Xuanhua Shi, Junling Liang, Sheng Di, Bingsheng He, Hai Jin 0001 |
IEEE Trans. Knowl. Data Eng. | 1 |
| 2018 | Deca: A Garbage Collection Optimizer for In-Memory Data ProcessingabstractIn-memory caching of intermediate data and active combining of data in shuffle buffers have been shown to be very effective in minimizing the recomputation and I/O cost in big data processing systems such as Spark and Flink. However, it has also been widely reported that these techniques would create a large amount of long-living data objects in the heap. These generated objects may quickly saturate the garbage collector, especially when handling a large dataset, and hence, limit the scalability of the system. To eliminate this problem, we propose a lifetime-based memory management framework, which, by automatically analyzing the user-defined functions and data types, obtains the expected lifetime of the data objects and then allocates and releases memory space accordingly to minimize the garbage collection overhead. In particular, we present Deca,1 a concrete implementation of our proposal on top of Spark, which transparently decomposes and groups objects with similar lifetimes into byte arrays and releases their space altogether when their lifetimes come to an end. When systems are processing very large data, Deca also provides field-oriented memory pages to ensure high compression efficiency. Extensive experimental studies using both synthetic and real datasets show that, in comparing to Spark, Deca is able to (1) reduce the garbage collection time by up to 99.9%, (2) reduce the memory consumption by up to 46.6% and the storage space by 23.4%, (3) achieve 1.2× to 22.7× speedup in terms of execution time in cases without data spilling and 16× to 41.6× speedup in cases with data spilling, and (4) provide similar performance compared to domain-specific systems. Xuanhua Shi, Zhixiang Ke, Yongluan Zhou, Hai Jin 0001, Lu Lu 0006, Ligang He |
ACM Trans. Comput. Syst. | 1 |
| 2018 | Core Maintenance in Dynamic Graphs: A Parallel Approach Based on MatchingabstractThe core number of vertices is a basic index depicting cohesiveness of a graph, and has been widely used in large-scale graph analytics. In this paper, we study the update of core numbers of vertices in dynamic graphs with edge insertions/deletions, which is known as the core maintenance problem. Different from previous approaches that just focus on the case of single-edge insertion/deletion and sequentially handle the edges when multiple edges are inserted/deleted, we investigate the parallelism in the core maintenance procedure. Specifically, we show that if the inserted/deleted edges constitute a matching, the core number update with respect to each inserted/deleted edge can be handled in parallel. Based on this key observation, we propose parallel algorithms for core maintenance in both cases of edge insertions and deletions. Extensive experiments are conducted to evaluate the efficiency, stability, parallelism and scalability of our algorithms on different types of real-world, synthetic graphs and temporal networks. Comparing with former approaches, our algorithms can improve the core maintenance efficiency significantly. Hai Jin 0001, Dongxiao Yu, Qiang-Sheng Hua, Xuanhua Shi |
IEEE Trans. Parallel Distributed Syst. | 5 |
| 2017 | TPS: An Efficient VM Scheduling Algorithm for HPC Applications in Cloud
Duoqiang Wang, Chi Zhang 0117, Xuanhua Shi, Hai Jin 0001 |
GPC | 4 |
| 2017 | Distributively Computing Random Walk Betweenness Centrality in Linear TimeabstractBetweenness centrality of a node represents its influence over the spread of information in the network. It is normally defined as the ratio of the number of shortest paths passing through the node among all shortest paths. However, the spread of information may not just pass through the shortest paths which is captured by a new measure of betweenness centrality based on random walks [1]. The random walk betweenness centrality of a node means how often it is traversed by a random walk between all pairs of other nodes. In this paper, we propose an O(n log n) time distributed randomized approximation algorithm for calculating each node's random walk betweenness centrality with an approximation ratio (1-ϵ) where n is the number of nodes and ϵ is an arbitrarily small constant between 0 and 1. Our distributed algorithm is designed under the widely used CONGEST model, where each edge can only transfer O(log n) bits in each round. To our best knowledge, this is the first distributed algorithm for computing the random walk betweenness centrality. Moreover, we give a non-trivial lower bound for distributively computing the exact random walk betweenness centrality under the CONGEST model, which is Ω(n\log n +D) where D is the network diameter. This means exactly computing random walk betweenness cannot be done in sublinear time. Qiang-Sheng Hua, Ming Ai, Hai Jin 0001, Dongxiao Yu, Xuanhua Shi |
ICDCS | 5 |
| 2017 | SSDUP: a traffic-aware ssd burst buffer for HPC systemsabstractMany high performance computing (HPC) applications are highly data intensive. Current HPC storage systems still use hard disk drives (HDDs) as their dominant storage devices, which suffer from disk head thrashing when accessing random data. New storage devices such as solid state drives (SSDs), which can handle random data access much more efficiently, have been widely deployed as the buffer to HDDs in many production HPC systems. Burst buffer has also been proposed to manage the SSD buffering of bursty write requests. Although burst buffer can improve I/O performance in many cases, we find that it has some limitations such as requiring large SSD capacity and harmonious overlapping between computation phase and data flushing stage. Xuanhua Shi, Wei Liu 0004, Hai Jin 0001, Chen Yu 0003, Yong Chen 0001 |
ICS | 1 |
| 2017 | MURS: Mitigating Memory Pressure in Service-Oriented Data Processing SystemabstractAlthough a data processing system often works as a batch processing system, many enterprises deploy such a system as a service, which we call the service-oriented data processing system. It has been shown that in-memory data processing systems suffer from serious memory pressure. The situation becomes even worse for the service-oriented data processing systems due to various reasons. For example, in a service-oriented system, multiple submitted tasks are launched at the same time and executed in the same context in the resources, compared with the batch processing mode where the tasks are processed one by one. Therefore, the memory pressure will affect all submitted tasks, including the tasks that only incur the light memory pressure when they are run alone. In this paper, we find that the reason why memory pressure arises is because the running tasks produce massive long-living data objects in the limited memory space. Our studies further reveal that the long-living data objects are generated by the API functions that are invoked by the in-memory processing frameworks. Based on these findings, we propose a method to classify the API functions based on the memory usage rate. Further, we design a scheduler called MURS to mitigate the memory pressure. We implement MURS in Spark and conduct the experiments to evaluate the performance of MURS. The results show that when comparing to Spark, MURS can 1) decrease the execution time of the submitted jobs by up to 65.8%, 2) mitigate the memory pressure in the server by decreasing the garbage collection time by up to 81%, and 3) reduce the data spilling, and hence disk I/O, by approximately 90%. Xuanhua Shi, Ligang He, Hai Jin 0001, Zhixiang Ke, Song Wu 0001 |
ICWS | 1 |
| 2017 | A task-based approach for finding SCCs in real-world graphs on external memoryabstractSummary Finding strongly connected components (SCCs) in graphs is one of the important research topics of graph data mining. Traditional SCC‐finding methods need to load the whole graph into main memory before actual processing, which makes them inappropriate to process today's large graphs. Although it is cost‐effective to conductive SCC‐finding in the external memory (EM) graph processing systems, existing EM systems are still inefficient on such workload for 2 reasons: data‐parallel processing model and inefficient graph mutation mechanism. In this paper, we propose a task‐based approach, named as TAS, for finding SCCs in large real‐world graphs. TAS encapsulates individual graph dataset and the SCC‐finding algorithms conducted on it as an independent task. By this strategy, different algorithms can be assigned to different datasets according to the stage of processing or the size of the dataset. By gradually removing unneeded data, ie, the known SCCs, with an efficient graph mutation method, TAS continuously reduces the scale of the problem to accelerate processing. Performance evaluation on large real‐world graphs show that TAS greatly outperforms existing EM solutions on finding SCCs (eg, 55.2x faster than GraphChi and 40.3x faster than X‐Stream). Huiming Lv, Zhiyuan Shao, Xuanhua Shi, Hai Jin 0001 |
Concurr. Comput. Pract. Exp. | 4 |
| 2017 | Optimizing Graph Processing on GPUsabstractDistributed vertex-centric model has been recently proposed for large-scale graph processing. Due to the simple but efficient programming abstraction, similar graph computing frameworks based on GPUs are gaining more and more attention. However, prior works of GPU-based graph processing suffer from load imbalance and irregular memory access because of the inherent characteristics of graph applications. In this paper, we propose a generalized graph computing framework for GPUs to simplify existing models but with higher performance. In particular, two novel algorithmic optimizations, lightweight approximate sorting and data layout transformation, are proposed to tackle the performance issues of current systems. With extensive experimental evaluation under a wide range of real world and synthetic workloads, we show that our system can achieve 1.6× to 4.5× speedups over the state-of-the-art. Wenyong Zhong, Jianhua Sun 0002, Hao Chen 0002, Zhiwen Chen 0006, Xuanhua Shi |
IEEE Trans. Parallel Distributed Syst. | 7 |
| 2016 | Measuring Directional Semantic Similarity with Multi-features
Bo Liu 0057, Xuanhua Shi, Hai Jin 0001 |
APWeb (1) | 2 |
| 2016 | Exploiting Sample Diversity in Distributed Machine Learning SystemsabstractWith the increase of machine learning scalability, there is a growing need for distributed systems which can execute machine learning algorithms on large clusters. Currently, most distributed machine learning systems are developed based on iterative optimization algorithm and parameter server framework. However, most systems compute on all samples in every iteration and this method consumes too much computing resources since the amount of samples is always too large. In this paper, we study on the sample diversity and find that most samples ontribute little to model updating during most iterations. Based on these findings, we propose a new iterative optimization algorithm to reduce the computation load by reusing the iterative computing results. The experiment demonstrates that, compared to the current methods, the algorithm proposed in this paper can reduce about 23% of the whole computation load without increasing of communications. Xuanhua Shi, Hai Jin 0001 |
CCGrid | 2 |
| 2016 | SSDUP: An Efficient SSD Write Buffer Using PipelineabstractHigh performance computing (HPC) applications are becoming more data-intensive and produce increasingly large I/O demands on storage systems. New storage devices such as SSD which has nearly no seek latency and high throughput have been widely used together with HDD to serve as a hybrid storage system. To solve the I/O bottleneck problem, existing hybrid storage solutions such as Burst Buffer have been proposed as intermediate layer between clients and disks to absorb burst I/O requests and improve write performance. However Burst Buffer needs sufficient SSD space to meet the maximum burst I/O requests which is still a costly solution. In this paper, we propose a hybrid architecture called SSDUP (an SSD write buffer Using Pipeline) for HPC storage systems, which uses NAND flash based SSD as a write-back buffer for HDD. With our efforts, SSDUP can achieve a good performance by using limited SSD space. Xuanhua Shi, Wei Liu 0004, Hai Jin 0001, Yong Chen 0001 |
CLUSTER | 2 |
| 2016 | Nearly Optimal Distributed Algorithm for Computing Betweenness CentralityabstractIn this paper, we propose an O(N) time distributed algorithm for computing betweenness centralities of all nodes in the network where N is the number of nodes. Our distributed algorithm is designed under the widely employed CONGEST model in the distributed computing community which limits each message only contains O(log N) bits. To our best knowledge, this is the first linear time deterministic distributed algorithm for computing the betweenness centralities in the published literature. We also give a lower bound for distributively computing the betweenness centrality under the CONGEST model as Ω(D+N/ log N) where D is the diameter of the network. This implies that our distributed algorithm is nearly optimal. Qiang-Sheng Hua, Haoqiang Fan, Ming Ai, Lixiang Qian, Xuanhua Shi, Hai Jin 0001 |
ICDCS | 6 |
| 2016 | Finding SCCs in Real-World Graphs on External Memory: A Task-Based ApproachabstractFinding Strongly Connected Components (SCCs) in graphs is one of the important research topics of graph data mining. Traditional methods of finding SCCs need to fully load the whole graph into the main memory of a computer before actual processing. However, with the rapid growth of real-world graphs, the sizes of graphs easily exceed the main memory space of an ordinary computer. The distributed graph processing system running on a cluster and the out-of-core system utilizing the external memory all can handle that huge graph, but recent evidences (e.g., GridGraph) show that the external memory systems are more cost-effective and efficient than the distributed systems on conducting most graph mining tasks. Existing external memory solutions are inefficient on finding SCCs in large-scale graphs for two reasons: (1) The data-parallel processing model adopted is not efficient to find SCCs in a largescale graph. (2) Their poor support for graph mutation incurs excessively high overhead. In this paper, we study the problem of finding SCCs in big real-world graphs by using the external memory. We propose a task-based approach and an efficient graph mutation method to address the limitations in existing external solutions for finding SCCs. Experiment results show that our approach is orders of magnitude faster than existing external memory solutions. Huiming Lv, Zhiyuan Shao, Xuanhua Shi, Hai Jin 0001 |
ISPDC | 3 |
| 2016 | Conch: A Cyclic MapReduce Model for Iterative ApplicationsabstractMapReduce programming model is a popular model to simplify but speed up data parallel applications. However, it is not efficient for iterative applications because of its repeated data transmission with HDFS (Hadoop Distributed File System). Conch, a cyclic MapReduce model, is designed for efficient processing of iterative applications. In order to minimize network overhead, shared data is cached locally and a "map-shuffle" phase is presented with a combined transmission mechanism. Meanwhile, a prediction scheduler for iterative applications is brought out to achieve better data locality in terms of runtime information. The experiments show that Conch can support iterative applications transparently and efficiently. Compared with Hadoop and HaLoop in single-job environment, Conch can achieve 13%-17% improvements on K-Means and fuzzy C-Means. Especially in multi-job environment, 63.6% and 28.6% improvements can be obtained compared with Hadoop and HaLoop. Genmao Yu, Hai Jin 0001, Xuanhua Shi, Qin Zhang 0004 |
PDP | 4 |
| 2016 | Brief Announcement: A Tight Distributed Algorithm for All Pairs Shortest Paths and ApplicationsabstractGiven an unweighted and undirected graph, this paper aims to give a tight distributed algorithm for computing the all pairs shortest paths (APSP) under synchronous communications and the CONGEST(B) model, where each node can only transfer B bits of information along each incident edge in a round. The best previous results for distributively computing APSP need O(N+D) time where N is the number of nodes and D is the diameter [1,2]. However, there is still a B factor gap from the lower bound Ω(N/B+D) [1]. In order to close this gap, we propose a multiplexing technique to push the parallelization of distributed BFS tree constructions to the limit such that we can solve APSP in O(N/B+D) time which meets the lower bound. This result also implies a Θ(N/B+D) time distributed algorithm for diameter. In addition, we extend our distributed algorithm to compute girth which is the length of the shortest cycle and clustering coefficient (CC) which is related to counting the number of triangles incident to each node. The time complexities for computing these two graph properties are also O(N/B+D). Qiang-Sheng Hua, Haoqiang Fan, Lixiang Qian, Ming Ai, Xuanhua Shi, Hai Jin 0001 |
SPAA | 6 |
| 2016 | Runtime-aware adaptive scheduling in stream processingabstractSummary Long‐running stream applications usually share the same fundamental computational infrastructure. To improve the efficiency of data processing in stream processing systems, a data analysis operator could be partitioned intonparallel tasks. The partitioned tasks are usually deployed onmnodes coexisting with other application operators. Because the node performance can vary in unpredictable ways (i.e., (1) stream input rates may fluctuate and (2) computational resource availability varies as other applications are affected), the nodes have different processing steps, and the slow node determines the operator performance. Hence, the tasks should be redistributed at runtime for stream applications to meet their strict latency requirements. Our key idea is to redistribute the tasks to the best node dynamically adaptive to resource or load fluctuations. In this paper, we present a runtime‐aware adaptive schedule mechanism that aims at minimizing the operator processing latency and minimizing the latency difference between different nodes' tasks. We propose a new abstraction calledperformance cost ratio(PCR) that evaluates the node performance. The higher the node's PCR is, the less cost the node will pay for processing one tuple, and the more tasks should be deployed on it. In a scheduling, we first sort tasks descendingly by their loads and sort nodes by their PCR. Then we reassign the amount of computation according to the node's PCR to keep the node's PCR and its input rate the same or in similar proportion in all PCRs. The PCR‐based quantitative algorithm applies itself to make tasks loads quantized to the processing capacity of nodes, move the minimum amount of operator's tasks, and keep the tasks local at the same time. We have implemented a runtime‐aware adaptive scheduler as an extension to Storm and evaluated this strategy. We achieve the optimization goal using less computational resources. Copyright © 2015 John Wiley & Sons, Ltd. Yuan Liu 0005, Xuanhua Shi, Hai Jin 0001 |
Concurr. Comput. Pract. Exp. | 2 |
| 2016 | Lifetime-Based Memory Management for Distributed Data Processing SystemsabstractIn-memory caching of intermediate data and eager combining of data in shuffle buffers have been shown to be very effective in minimizing the re-computation and I/O cost in distributed data processing systems like Spark and Flink. However, it has also been widely reported that these techniques would create a large amount of long-living data objects in the heap, which may quickly saturate the garbage collector, especially when handling a large dataset, and hence would limit the scalability of the system. To eliminate this problem, we propose a lifetime-based memory management framework, which, by automatically analyzing the user-defined functions and data types, obtains the expected lifetime of the data objects, and then allocates and releases memory space accordingly to minimize the garbage collection overhead. In particular, we present Deca, a concrete implementation of our proposal on top of Spark, which transparently decomposes and groups objects with similar lifetimes into byte arrays and releases their space altogether when their lifetimes come to an end. An extensive experimental study using both synthetic and real datasets shows that, in comparing to Spark, Deca is able to 1) reduce the garbage collection time by up to 99.9%, 2) to achieve up to 22.7x speed up in terms of execution time in cases without data spilling and 41.6x speedup in cases with data spilling, and 3) to consume up to 46.6% less memory. Lu Lu 0006, Xuanhua Shi, Yongluan Zhou, Hai Jin 0001, Cheng Pei, Ligang He, Yuanzhen Geng |
Proc. VLDB Endow. | 2 |
| 2016 | FITDOC: fast virtual machines checkpointing with delta memory compression
Yunjie Du, Xuanhua Shi, Hai Jin 0001, Song Wu 0001, Laurence T. Yang |
J. Supercomput. | 2 |
| 2016 | Poris: A Scheduler for Parallel Soft Real-Time Applications in Virtualized EnvironmentsabstractWith the prevalence of cloud computing and virtualization, more and more cloud services including parallel soft real-time applications (PSRT applications) are running in virtualized data centers. However, current hypervisors do not provide adequate support for them because of soft real-time constraints and synchronization problems, which result in frequent deadline misses and serious performance degradation. CPU schedulers in underlying hypervisors are central to these issues. In this paper, we identify and analyze CPU scheduling problems in hypervisors. Then, we design and implement a parallel soft real-time scheduler according to the analysis, namedPoris, based on Xen. It addresses both soft real-time constraints and synchronization problems simultaneously. In our proposed method,priority promotionanddynamic time slicemechanisms are introduced to determine when to schedulevirtual CPUs(VCPUs) according to the characteristics of soft real-time applications. Besides, considering that PSRT applications may run in avirtual machine(VM) or multiple VMs, we presentparallel scheduling,group schedulingandcommunication-driven group schedulingto accelerate synchronizations of these applications and make sure that tasks are finished before their deadlines under different scenarios. Our evaluation showsPoriscan significantly improve the performance of PSRT applications no matter how they run in a VM or multiple VMs. For example, compared to the Credit scheduler,Porisdecreases the response time of web search benchmark by up to 91.6 percent. Song Wu 0001, Like Zhou, Huahua Sun, Hai Jin 0001, Xuanhua Shi |
IEEE Trans. Parallel Distributed Syst. | 5 |
| 2015 | Towards Truly Elastic Distributed Graph Computing in the Cloud
Lu Lu 0006, Xuanhua Shi, Hai Jin 0001 |
APSCC | 2 |
| 2015 | Improving the Memory Efficiency of In-Memory MapReduce Based HPC Systems
Cheng Pei, Xuanhua Shi, Hai Jin 0001 |
ICA3PP (1) | 2 |
| 2015 | Optimization of asynchronous graph processing on GPU with hybrid coloring modelabstractModern GPUs have been widely used to accelerate the graph processing for complicated computational problems regarding graph theory. Many parallel graph algorithms adopt the asynchronous computing model to accelerate the iterative convergence. Unfortunately, the consistent asynchronous computing requires locking or the atomic operations, leading to significant penalties/overheads when implemented on GPUs. To this end, coloring algorithm is adopted to separate the vertices with potential updating conflicts, guaranteeing the consistency/correctness of the parallel processing. We propose a light-weight asynchronous processing framework called Frog with a hybrid coloring model. We find that majority of vertices (about 80%) are colored with only a few colors, such that they can be read and updated in a very high degree of parallelism without violating the sequential consistency. Accordingly, our solution will separate the processing of the vertices based on the distribution of colors. Xuanhua Shi, Junling Liang, Sheng Di, Bingsheng He, Hai Jin 0001, Lu Lu 0006, Jianlong Zhong |
PPoPP | 1 |
| 2015 | Towards Optimized Fine-Grained Pricing of IaaS Cloud PlatformabstractAlthough many pricing schemes in IaaS platform are already proposed with pay-as-you-go and subscription/spot market policy to guarantee service level agreement, it is still inevitable to suffer from wasteful payment because of coarse-grained pricing scheme. In this paper, we investigate an optimized fine-grained and fair pricing scheme. Two tough issues are addressed: (1) the profits of resource providers and customers often contradict mutually; (2) VM-maintenance overhead like startup cost is often too huge to be neglected. Not only can we derive an optimal price in the acceptable price range that satisfies both customers and providers simultaneously, but we also find a best-fit billing cycle to maximize social welfare (i.e., the sum of the cost reductions for all customers and the revenue gained by the provider). We carefully evaluate the proposed optimized fine-grained pricing scheme with two large-scale real-world production traces (one from Grid Workload Archive and the other from Google data center). We compare the new scheme to classic coarse-grained hourly pricing scheme in experiments and find that customers and providers can both benefit from our new approach. The maximum social welfare can be increased up to 72.98 and 48.15 percent with respect to DAS-2 trace and Google trace respectively. Hai Jin 0001, Xinhou Wang, Song Wu 0001, Sheng Di, Xuanhua Shi |
IEEE Trans. Cloud Comput. | 5 |
| 2015 | Mammoth: Gearing Hadoop Towards Memory-Intensive MapReduce ApplicationsabstractThe MapReduce platform has been widely used for large-scale data processing and analysis recently. It works well if the hardware of a cluster is well configured. However, our survey has indicated that common hardware configurations in small- and medium-size enterprises may not be suitable for such tasks. This situation is more challenging for memory-constrained systems, in which the memory is a bottleneck resource compared with the CPU power and thus does not meet the needs of large-scale data processing. The traditional high performance computing (HPC) system is an example of the memory-constrained system according to our survey. In this paper, we have developed Mammoth, a new MapReduce system, which aims to improve MapReduce performance using global memory management. In Mammoth, we design a novel rule-based heuristic to prioritize memory allocation and revocation among execution units (mapper, shuffler, reducer, etc.), to maximize the holistic benefits of the Map/Reduce job when scheduling each memory unit. We have also developed a multi-threaded execution engine, which is based on Hadoop but runs in a single JVM on a node. In the execution engine, we have implemented the algorithm of memory scheduling to realize global memory management, based on which we further developed the techniques such as sequential disk accessing, multi-cache and shuffling from memory, and solved the problem of full garbage collection in the JVM. We have conducted extensive experiments to compare Mammoth against the native Hadoop platform. The results show that the Mammoth system can reduce the job execution time by more than 40 percent in typical cases, without requiring any modifications of the Hadoop programs. When a system is short of memory, Mammoth can improve the performance by up to 5.19 times, as observed for I/O intensive applications, such as PageRank. We also compared Mammoth with Spark. Although Spark can achieve better performance than Mammoth for interactive and iterative applications when the memory is sufficient, our experimental results show that for batch processing applications, Mammoth can adapt better to various memory environments and outperform Spark when the memory is insufficient, and can obtain similar performance as Spark when the memory is sufficient. Given the growing importance of supporting large-scale data processing and analysis and the proven success of the MapReduce platform, the Mammoth system can have a promising potential and impact. Xuanhua Shi, Ligang He, Lu Lu 0006, Hai Jin 0001, Yong Chen 0001, Song Wu 0001 |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2015 | Synchronization-Aware Scheduling for Virtual Clusters in CloudabstractDue to high flexibility and cost-effectiveness, cloud computing is increasingly being explored as an alternative to local clusters by academic and commercial users. Recent research already confirmed the feasibility of running tightly-coupled parallel applications with virtual clusters. However, such types of applications suffer from significant performance degradation, especially as the over-commitment is common in cloud. That is, the number of executable Virtual CPUs (VCPUs) is often larger than that of available Physical CPUs (PCPUs) in the system. The performance degradation is mainly due to the fact that the current virtual machine monitors (VMMs) are unaware of the synchronization requirements of the VMs which are running parallel applications. In this paper, There are two key contributions. (1) We propose an autonomous synchronization-aware VM scheduling (SVS) algorithm, which can effectively mitigate the performance degradation of tightly-coupled parallel applications running atop them in over-committed situation. (2) We integrate the SVS algorithm into Xen VMM scheduler, and rigorously implement a prototype. We evaluate our design on a real cluster environment with NPB benchmark and real-world trace. Experiments show that our solution attains better performance for tightly-coupled parallel applications than the state-of-the-art approaches like Xen's Credit scheduler, balance scheduling, and hybrid scheduling. Song Wu 0001, Haibao Chen, Sheng Di, Bing Bing Zhou, Zhenjiang Xie, Hai Jin 0001, Xuanhua Shi |
IEEE Trans. Parallel Distributed Syst. | 7 |
| 2014 | Iteration Based Collective I/O Strategy for Parallel I/O SystemsabstractMPI collective I/O is a widely used I/O method that helps data-intensive scientific applications gain better I/O performance. However, it has been observed that existing collective I/O strategies do not perform well due to the access contention problem. Existing collective I/O optimization strategies mainly focus on the I/O phase efficiency and ignore the shuffle cost that may limit the potential of their performance improvement. We observe that as the size of I/O becomes larger, one I/O operation from the upper application would be separated into several iterations to complete. So, I/O requests in each file domain do not necessarily issue to the parallel file system simultaneously unless they are carried out within the same iteration step. Based on that observation, this paper proposes a new collective I/O strategy that reorganizes I/O requests within each file domain instead of coordinating requests across file domains, such that we can eliminate access contentions without introducing extra shuffle cost between aggregators and computing processes. Using benchmark workloads IOR, we evaluate our new strategy and compare with the conventional one. The proposed strategy achieves up to 47%-63% I/O bandwidth improvement compared to the existing ROMIO collective I/O strategy. Xuanhua Shi, Hai Jin 0001, Song Wu 0001, Yong Chen 0001 |
CCGRID | 2 |
| 2014 | An adaptive PSO based on motivation mechanism and acceleration restraint operator
Jiangshao Gu, Xuanhua Shi |
IEEE Congress on Evolutionary Computation | 2 |
| 2014 | GIRAFFE: A scalable distributed coordination service for large-scale systemsabstractThe scale of cloud services keeps increasing over time, significantly introducing huge challenges in system manageability and reliability. Designing coordination services in cloud is the right track to solve the above problems. However, existing coordination services (e.g., Chubby and ZooKeeper) only perform well in read-intensive scenario and small ensemble scales. To this end, we propose Giraffe, a scalable distributed coordination service. There are three important contributions in our design. (1) Giraffe organizes coordination servers using interior-node-disjoint trees for better scalability. (2) Giraffe employs a novel Paxos protocol for strong consistency and fault-tolerance. (3) Giraffe supports hierarchical data organization and in-memory storage for high throughput and low latency. We evaluate Giraffe on a high performance computing test-bed. The experimental results show that Giraffe gains much better write performance than ZooKeeper when server ensemble is large. Giraffe is nearly 300% faster than ZooKeeper on update operations when ensemble size is 50 servers. Experiments also show that Giraffe reacts and recovers more quickly than ZooKeeper against node failures. Xuanhua Shi, Haohong Lin, Hai Jin 0001, Bing Bing Zhou, Zuoning Yin, Sheng Di, Song Wu 0001 |
CLUSTER | 1 |
| 2014 | Communication-driven scheduling for virtual clusters in cloudabstractDue to high flexibility and cost-effectiveness, cloud computing is increasingly being explored as an alternative to local clusters by academic and commercial users. Recent research already confirmed the feasibility of running tightly-coupled parallel applications with virtual clusters. However, such types of applications suffer from significant performance degradation, especially as the overcommitment is common in cloud. That is, the number of executable Virtual CPUs (VCPUs) is often larger than that of available Physical CPUs (PCPUs) in the system. The performance degradation mainly results from that the current Virtual Machine Monitors (VMMs) cannot co-schedule (or coordinate at the same time) the VCPUs that host parallel application threads/processes with synchronization requirements. Haibao Chen, Song Wu 0001, Sheng Di, Bing Bing Zhou, Zhenjiang Xie, Hai Jin 0001, Xuanhua Shi |
HPDC | 7 |
| 2014 | A Real-Time Scheduling Framework Based on Multi-core Dynamic Partitioning in Virtualized Environment
Song Wu 0001, Like Zhou, Danqing Fu, Hai Jin 0001, Xuanhua Shi |
NPC | 5 |
| 2014 | Page Classifier and Placer: A Scheme of Managing Hybrid Caches
Xuanhua Shi, Hai Jin 0001, Xiaofei Liao, Song Wu 0001, Xiaoming Li 0010 |
NPC | 2 |
| 2014 | MECOM: Live migration of virtual machines by adaptively compressing memory pages
Hai Jin 0001, Song Wu 0001, Xuanhua Shi, Hanhua Chen |
Future Gener. Comput. Syst. | 4 |
| 2014 | Morpho: A decoupled MapReduce framework for elastic cloud computing
Lu Lu 0006, Xuanhua Shi, Hai Jin 0001, Qiuyue Wang, Daxing Yuan, Song Wu 0001 |
Future Gener. Comput. Syst. | 2 |
| 2014 | Dynamic and fast processing of queries on large-scale RDF data
Pingpeng Yuan, Changfeng Xie, Hai Jin 0001, Ling Liu 0001, Xuanhua Shi |
Knowl. Inf. Syst. | 6 |
| 2013 | Cost-Aware Client-Side File Caching for Data-Intensive ApplicationsabstractParallel and distributed file systems are widely used to provide high throughput in high-performance computing and Cloud computing systems. To increase the parallelism, I/O requests are partitioned into multiple sub-requests (or `flows') and distributed across different data nodes. The performance of file systems is extremely poor if data nodes have highly unbalanced response time. Client-side caching offers a promising direction for addressing this issue. However, current work has primarily used client-side memory as a read cache and employed a write-through policy which requires synchronous update for every write and significantly under-utilizes the client-side cache when the applications are write-intensive. Realizing that the cost of an I/O request depends on the struggler sub-requests, we propose a cost-aware client-side file caching (CCFC) strategy, that is designed to cache the sub-requests with high I/O cost on the client end. This caching policy enables a new trade-off across write performance, consistency guarantee and cache size dimensions. Using benchmark workloads MADbench2, we evaluate our new cache policy alongside conventional write-through. We find that the proposed CCFC strategy can achieve up to 110% throughput improvement compared to the conventional write-through policies with the same cache size on an 85-node cluster. Yaning Huang, Hai Jin 0001, Xuanhua Shi, Song Wu 0001, Yong Chen 0001 |
CloudCom (2) | 3 |
| 2013 | Supporting parallel soft real-time applications in virtualized environment
Like Zhou, Song Wu 0001, Huahua Sun, Hai Jin 0001, Xuanhua Shi |
HPDC | 5 |
| 2013 | Virtual Machine Scheduling for Parallel Soft Real-Time ApplicationsabstractWith the prevalence of multicore processors in computer systems, many soft real-time applications, such as media-based ones, use parallel programming models to utilize hardware resources better and possibly shorten response time. Meanwhile, virtualization technology is widely used in cloud data centers. More and more cloud services including such parallel soft real-time applications are running in virtualized environment. However, current hyper visors do not provide adequate support for them because of soft real-time constraints and synchronization problems, which result in frequent deadline misses and serious performance degradation. CPU schedulers in underlying hyper visors are central to these issues. In this paper, we identify and analyze CPU scheduling problems in hyper visors, and propose a novel scheduling algorithm considering both soft real-time constraints and synchronization problems. In our proposed method, real-time priority is introduced to accelerate event processing of parallel soft real-time applications, and dynamic time slice is used to schedule virtual CPUs. Besides, all runnable virtual CPUs of virtual machines running parallel soft real-time applications are scheduled simultaneously to address synchronization problems. We implement a parallel soft real-time scheduler, named Poris, based on Xen. Our evaluation shows Poris can significantly improve the performance of parallel soft real-time applications. For example, compared to the Credit scheduler, Poris improves the performance of media player by up to a factor of 1.35, and shortens the execution time of PARSEC benchmark by up to 44.12%. Like Zhou, Song Wu 0001, Huahua Sun, Hai Jin 0001, Xuanhua Shi |
MASCOTS | 5 |
| 2013 | Dependable Grid Workflow Scheduling Based on Resource Availability
Yongcai Tao, Hai Jin 0001, Song Wu 0001, Xuanhua Shi, Lei Shi 0001 |
J. Grid Comput. | 4 |
| 2013 | Developing an optimized application hosting framework in Clouds
Xuanhua Shi, Hongbo Jiang 0001, Ligang He, Hai Jin 0001, Chonggang Wang, Xueguang Chen |
J. Comput. Syst. Sci. | 1 |
| 2013 | SafeStack: Automatically Patching Stack-Based Buffer Overflow VulnerabilitiesabstractBuffer overflow attacks still pose a significant threat to the security and availability of today's computer systems. Although there are a number of solutions proposed to provide adequate protection against buffer overflow attacks, most of existing solutions terminate the vulnerable program when the buffer overflow occurs, effectively rendering the program unavailable. The impact on availability is a serious problem on service-oriented platforms. This paper presents SafeStack, a system that can automatically diagnose and patch stack-based buffer overflow vulnerabilities. The key technique of our solution is to virtualize memory accesses and move the vulnerable buffer into protected memory regions, which provides a fundamental and effective protection against recurrence of the same attack without stopping normal system execution. We developed a prototype on a Linux system, and conducted extensive experiments to evaluate the effectiveness and performance of the system using a range of applications. Our experimental results showed that SafeStack can quickly generate runtime patches to successfully handle the attack's recurrence. Furthermore, SafeStack only incurs acceptable overhead for the patched applications. Hai Jin 0001, Deqing Zou, Bing Bing Zhou, Zhenkai Liang, Weide Zheng, Xuanhua Shi |
IEEE Trans. Dependable Secur. Comput. | 7 |
| 2013 | Adapting grid computing environments dependable with virtual machines: design, implementation, and evaluations
Xuanhua Shi, Hai Jin 0001, Song Wu 0001 |
J. Supercomput. | 1 |
| 2013 | A disk bandwidth allocation mechanism with priority
Xibin Wang, Hai Jin 0001, Xuanhua Shi, Wenzhi Cao, Xijiang Ke |
J. Supercomput. | 4 |
| 2012 | Effectively deploying services on virtualization infrastructure
Hai Jin 0001, Song Wu 0001, Xuanhua Shi, Jinyan Yuan |
Frontiers Comput. Sci. | 4 |
| 2012 | Toward scalable Web systems on multicore clusters: making use of virtual machines
Xuanhua Shi, Hai Jin 0001, Hongbo Jiang 0001, Dachuan Huang |
J. Supercomput. | 1 |
| 2011 | Optimizing the live migration of virtual machine by CPU scheduling
Hai Jin 0001, Song Wu 0001, Xuanhua Shi, Xiaoxin Wu 0001 |
J. Netw. Comput. Appl. | 4 |
| 2010 | MR-scope: a real-time tracing tool for MapReduceabstractMapReduce programming model is emerging as an efficient tool for data-intensive applications. Hadoop, an open-source implementation of MapReduce, has been widely adopted and experienced by both academia and enterprise. Recently, lots of efforts have been done on improving the performance of MapReduce system and on analyzing the MapReduce process based on the log files generated during the Hadoop execution. Visualizing log files seems to be a very useful tool to understand the behavior of the Hadoop process. In this paper, we present MR-Scope, a real-time MapReduce tracing tool. MR-Scope provides a real-time insight of the MapReduce process, including the ongoing progress of every task hosted in Task Tracker. In addition, it displays the health of the Hadoop cluster data nodes, the distribution of the file system's blocks and their replicas and the content of the different block splits of the file system. We implement MR-Scope in native Hadoop 0.1. Experimental results demonstrat that MR-Scope's overhead is less than 4% when running wordcount benchmark. Dachuan Huang, Xuanhua Shi, Shadi Ibrahim, Lu Lu 0006, Hongzhang Liu, Song Wu 0001, Hai Jin 0001 |
HPDC | 2 |
| 2010 | VirtCFT: A Transparent VM-Level Fault-Tolerant System for Virtual ClustersabstractA virtual cluster consists of a multitude of virtual machines and software components that are doomed to fail eventually. In many environments, such failures can result in unanticipated, potentially devastating failure behavior and in service unavailability. The ability of failover is essential to the virtual cluster's availability, reliability, and manageability. Most of the existing methods have several common disadvantages: requiring modifications to the target processes or their OSes, which is usually error prone and sometimes impractical; only targeting at taking checkpoints of processes, not whole entire OS images, which limits the areas to be applied. In this paper we present VirtCFT, an innovative and practical system of fault tolerance for virtual cluster. VirtCFT is a system-level, coordinated distributed checkpointing fault tolerant system. It coordinates the distributed VMs to periodically reach the globally consistent state and take the checkpoint of the whole virtual cluster including states of CPU, memory, disk of each VM as well as the network communications. When faults occur, VirtCFT will automatically recover the entire virtual cluster to the correct state within a few seconds and keep it running. Superior to all the existing fault tolerance mechanisms, VirtCFT provides a simpler and totally transparent fault tolerant platform that allows existing, unmodified software and operating system (version unawareness) to be protected from the failure of the physical machine on which it runs. We have implemented this system based on the Xen virtualization platform. Our experiments with real-world benchmarks demonstrate the effectiveness and correctness of VirtCFT. Minjia Zhang, Hai Jin 0001, Xuanhua Shi, Song Wu 0001 |
ICPADS | 3 |
| 2010 | Virtual Machine Management Based on Agent ServiceabstractWith the popularity of virtualization, the problem that how to manage hundreds even thousands of virtual machines running on multiple physical computing nodes becomes important. Current virtual machine management systems only can obtain basic information of virtual machines and execute simple operations on them, such as start, reboot and shutdown. In this paper, we design a virtual machine management approach based on agent service. Agent service can provide detail running status information inside virtual machines. It also has been a bridge for host machines and virtual machines to interact with each other. Agent service is designed to automatically start when virtual machine boots up. By agent service we can get real-time information about virtual machines. We evaluate monitoring overhead and the performance of batch operations when using agent. The experimental results show that agent mechanism outperforms methods using Libvirt or VMware tools. Song Wu 0001, Hai Jin 0001, Xuanhua Shi, Yankun Zhao, Jianyin Zhang |
PDCAT | 4 |
| 2010 | Adapting grid applications to safety using fault-tolerant methods: Design, implementation and evaluations
Xuanhua Shi, Jean-Louis Pazat, Eric Rodriguez, Hai Jin 0001, Hongbo Jiang 0001 |
Future Gener. Comput. Syst. | 1 |
| 2010 | Scalable DHT- and ontology-based information service for large-scale grids
Yongcai Tao, Hai Jin 0001, Song Wu 0001, Xuanhua Shi |
Future Gener. Comput. Syst. | 4 |
| 2010 | DAGMap: efficient and dependable scheduling of DAG workflow job in Grid
Haijun Cao, Hai Jin 0001, Xiaoxin Wu 0001, Song Wu 0001, Xuanhua Shi |
J. Supercomput. | 5 |
| 2009 | Evaluating MapReduce on Virtual Machines: The Hadoop Case
Shadi Ibrahim, Hai Jin 0001, Lu Lu 0006, Song Wu 0001, Xuanhua Shi |
CloudCom | 6 |
| 2009 | Live virtual machine migration with adaptive, memory compressionabstractLive migration of virtual machines has been a powerful tool to facilitate system maintenance, load balancing, fault tolerance, and power-saving, especially in clusters or data centers. Although pre-copy is a predominantly used approach in the state of the art, it is difficult to provide quick migration with low network overhead, due to a great amount of transferred data during migration, leading to large performance degradation of virtual machine services. This paper presents the design and implementation of a novel memory-compression-based VM migration approach (MECOM) that first uses memory compression to provide fast, stable virtual machine migration, while guaranteeing the virtual machine services to be slightly affected. Based on memory page characteristics, we design an adaptive zero-aware compression algorithm for balancing the performance and the cost of virtual machine migration. Pages are quickly compressed in batches on the source and exactly recovered on the target. Experiment demonstrates that compared with Xen, our system can significantly reduce 27.1% of downtime, 32% of total migration time and 68.8% of total transferred data on average. Hai Jin 0001, Song Wu 0001, Xuanhua Shi |
CLUSTER | 4 |
| 2008 | WAGA: A Flexible Web-Based Framework for Grid ApplicationsabstractThe research about the interaction between grid environments and users is becoming popular. Many scientists propose the integration for web 2.0 technologies and grid computing. In this paper, we propose a flexible web-based framework for grid applications - WAGA, which tries to bridge the gap between grid middleware and grid applications. WAGA provides a WYSIWYG (what you see is what you get) way for the programming for grid applications by adopting participation, interaction and sharing features of web 2.0 technology. With WAGA, a grid user is able to use the grid resources and to develop grid applications without understanding the underlying complexity of grids. WAGA is composed with a web GUI (called WAGA-designer) and some web-based APIs. WAGA-designer is used by grid users to develop application-based web portal, and the webAPIs are used by the WAGA-designer. The use case study shows that WAGA is flexible for grid users, and the performance evaluation shows that WAGA works with high efficiency. Xuanhua Shi, Hai Jin 0001, Song Wu 0001 |
APSCC | 1 |
| 2008 | Effectively Deploying Virtual Machines on ClusterabstractVirtualization technology has provided an opportunity to the efficient usage of computing resources. However, the management of VMs on cluster is still in the preliminary stage. How to construct user’s task environments fastly and efficiently remains a significant challenge. This paper presents a Multiple-VM Deployment System (MVDS)for creating and configuring users’ task environments on-demand. The system provides a template management model and all the VMs are created based on the templates including operating systems and applications. To improve the deployment performance, we explore some strategies about incremental mechanism and deployment tactics. We evaluate both the deployment time and I/O performance with proposed incremental mechanism. The experimental results show that the incremental mechanism outperforms the clone tactic. Song Wu 0001, Jinyan Yuan, Xuanhua Shi, Hai Jin 0001 |
APSCC | 4 |
| 2008 | ADVE: Adaptive and Dependable Virtual Environments for Grid Computing
Xuanhua Shi, Hai Jin 0001 |
GPC | 1 |
| 2008 | Dynasa: adapting grid applications to safety using fault-tolerant methodsabstractGrid applications have been prone to encountering problems such as failures or malicious attacks during execution, due to their distributed and large-scale features. The application itself, however, has limited power to address these problems. This paper presents the design, and implementation of an adaptive framework - Dynasa, which strives to handle security problems using adaptive fault-tolerance (i.e., checkpointing and replication) during the execution of applications according to the status of the grid environments. Xuanhua Shi, Jean-Louis Pazat, Eric Rodriguez, Hai Jin 0001, Hongbo Jiang 0001 |
HPDC | 1 |
| 2007 | Service Dependency Model for Dynamic and Stateful Grid Services
Hai Jin 0001, Yaqin Luo, Xuanhua Shi, Chengwei Wang |
ICA3PP | 4 |
| 2006 | Peer-Tree: A Hybrid Peer-to-Peer Overlay for Service DiscoveryabstractEfficient service discovery in dynamic, crossorganizational is one of the challenge aspects in ChinaGrid. Network overlay and search algorithms are two important considerations to the problem. Tree topology of organizations is easily managed but the root is a single point of failure; P2P structure tends to be conversed. To merge the advantages of both, we present a hybrid system with two layers: Tree Layer and Peer Layer. This structure is practical because organizations targeting at sub objects of a subject are inclined to be organized hierarchically as a Superpeer; and all these Superpeers construct an unstructured P2P network that adapts to peer’s interest by Learn Neighbor Algorithm. Experimental evaluation shows that our mechanism exhibits better search performance than fully hierarchical or fully distributed system. Jing Tie, Hai Jin 0001, Shengli Li 0002, Xuanhua Shi, Hanhua Chen, Xiaoming Ning |
AINA (1) | 4 |
| 2005 | VO-Sec: An Access Control Framework for Dynamic Virtual Organization
Hai Jin 0001, Weizhong Qiang, Xuanhua Shi, Deqing Zou |
ACISP | 3 |
| 2005 | A Formal General Framework and Service Access Model for Service GridabstractConstituent resources in a grid system need to be used in a coordinated fashion to deliver non trivial qualities of service. Various e-science and e-business use cases are investigated to guide how to create grid systems, and determine which functions grid systems should have. Web services emerge as a standard interoperable technology for grid systems. Although the motivations and goals for service grids are obvious, there is no clear definition for service grids to define and describe the general framework and service access model. In this paper, the general framework for service grids is defined in a formal approach, and the virtual organization based service access mechanism is modeled based on abstract state machines (ASM). In the service access model we proposed, the quality of service (QoS) issue is considered for the user request. This resulting serves as a theoretical base for our service grid system, HowU. Deqing Zou, Weizhong Qiang, Xuanhua Shi |
ICECCS | 3 |
| 2005 | A Virtual-Service-Domain Based Bidding Algorithm for Resource Discovery in Computational GridabstractResource discovery is a basic service in grid computing: gives a description of resources desired and finds the available one to match the description. In computational grid, how to discover resources efficiently has become a crucial factor to evaluate the performance in the whole system. In this paper, we present a bid-based resource discovery algorithm, which converts a resource request into a bidding letter and sends it to a group of physical services owned by the same virtual service to call for bidding. All resources receiving bidding letter make offers to bid according to our algorithm. Job manager selects the best one to response client request. To evaluate the performance of our method, we compare our system with the centralized and peer-to-peer resource discovery approaches. The analysis results show that our system reduces average response time of jobs, leverages the cost of the resource discovery, and improves the system scalability. Hongbo Zou, Hai Jin 0001, Zongfen Han, Jing Tie, Xuanhua Shi |
Web Intelligence | 5 |