EDBT 2026 Demo / reviewers in the wild / expert
Dong Dai 0001
dblp:55/6260-1
· DBLP profile ↗
61ranked-venue papers
15as first author
28since 2021 · last 2026
0000-0003-4078-8149ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 53 · 14 first-author · 26 since 2021Applied, interdisciplinary, general and emerging computing · 4 · 1 first-author · 1 since 2021Artificial intelligence and machine learning · 3 · 1 first-author · 1 since 2021Software engineering, systems software and programming languages · 2Databases, data management, data science and information retrieval · 2 · 1 first-author · 1 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | QoSFlow: Ensuring Service Quality of Distributed Workflows Using Interpretable Sensitivity Models
Md. Hasanur Rashid, Jesun Sahariar Firoz, Nathan R. Tallent, Luanzheng Guo, Dong Dai 0001 |
IPDPS | 6 |
| 2026 | CARAT: Client-Side Adaptive RPC and Cache Co-Tuning for Parallel File Systems
Md. Hasanur Rashid, Nathan R. Tallent, Forrest Sheng Bao, Dong Dai 0001 |
IPDPS | 4 |
| 2026 | TSUE+: An Efficient Update Framework With Swift Recycling Mechanism for Erasure-Coded Cluster File SystemsabstractCompared to replication-based storage systems, erasure-coded storage incurs significantly higher overhead during data updates. To address this issue, various parity logging methods have been proposed. Nevertheless, due to the long update path and substantial amount of random I/O involved in erasure code update processes, the resulting long latency and low through put often fail to meet the requirements of high performance applications. To address this challenge, we propose TSUE+, an efficient update framework with a swift recycling mechanism. TSUE+ divides the update process into two distinct stages: in the synchronous stage, data updates are stored in the format of replica data logs, eliminating random I/O by trading space for time; in the asynchronous stage, the recorded update logs are recycled and merged into original data and parity blocks, thereby reclaiming the storage overhead incurred in the synchronization phase. By converting random I/O operations into sequential ones based on data logs, TSUE+ effectively reduces update latency; furthermore, it significantly minimizes recycling overhead using a three-layer log structure and by leveraging the spatio-temporal locality of access patterns. We evaluated TSUE+ and other state of-the-art (SOTA) update mechanisms under diverse encoding schemes, using heterogeneous storage devices—including HDDs, SATA SSDs, NVMe SSDs, and PMEM—and multiple real-world and synthetic workloads: the MSR Cambridge trace, the Alibaba Cloud trace, the Tencent Cloud trace, and multiple synthetic worst-case workloads. Among all the platforms, TSUE+ has achieved significant performance improvements compared to other update methods, it also indicates that TSUE+ can be applied to various storage devices. Additionally, we provided percentile-based tail latency tests and update tests under the worst-case environment, which further demonstrated the ro bustness of TSUE+. Moreover, by enabling prompt log recycle and avoiding unnecessary overwrites and improving update granularity through locality-aware recycling, TSUE+ not only improves update performance but also mitigates write wear on SSD devices, thereby extending their operational lifespan. Yida Gu, Wenjing Huang 0002, Yili Ma, Dong Dai 0001, Guangming Tan, Dingwen Tao |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2025 | Dial: Decentralized I/O Autotuning Via Learned Client-Side Local Metrics for Parallel File SystemabstractEnabling efficient, high-performance data access in parallel file systems (PFS) is critical for today's highperformance computing systems. PFS client-side I/O heavily impacts the final I/O performance delivered to individual applications and the entire system. Autotuning the key client-side I/O behaviors has been extensively studied and shows promising results. However, existing work has heavily relied on extensive number of global runtime metrics to monitor and accurate modeling of applications' I/O patterns. Such heavy overheads significantly limit the ability to enable fine-grained, dynamic tuning in practical systems. In this study, we propose DIAL (Decentralized I/O AutoTuning via Learned Client-side Local Metrics) which takes a drastically different approach. Instead of trying to extract the global I/O patterns of applications, DIAL takes a decentralized approach, treating each I/O client as an independent unit and tuning configurations using only its locally observable metrics. With the help of machine learning models, DIAL enables multiple tunable units to make independent but collective decisions, reacting to what is happening in the global storage systems in a timely manner and achieving better I/O performance globally for the application. Md. Hasanur Rashid, Youbiao He, Forrest Sheng Bao, Dong Dai 0001 |
CCGrid | 5 |
| 2025 | TSUE: A Two-Stage Data Update Method for an Erasure Coded Cluster File SystemabstractCompared to replication-based storage systems, erasure-coded storage incurs significantly higher overhead during data updates. To address this issue, various parity logging methods have been proposed. Nevertheless, due to the long update path and substantial amount of random I/O involved in erasure code update processes, the resulting long latency and low throughput often fail to meet the requirements of high performance applications. To this end, we propose a two-stage data update method called TSUE. TSUE divides the update process into a synchronous stage that records updates in a data log, and an asynchronous stage that recycles the log in real-time. TSUE effectively reduces update latency by transforming random I/O into sequential I/O, and it significantly reduces recycle overhead by utilizing a three-layer log and the spatio-temporal locality of access patterns. In SSDs cluster, TSUE significantly improves update performance, achieving improvements of 7.6× under Ali-Cloud trace, 5× under Ten-Cloud trace, while it also extends the SSD's lifespan by up to 13× through reducing the frequencies of reads/writes and of erase operations. Yida Gu, Wenjing Huang 0002, Dong Dai 0001, Guangming Tan, Dingwen Tao |
HPDC | 5 |
| 2025 | IOAgent: Democratizing Trustworthy HPC I/O Performance Diagnosis Capability via LLMsabstractAs the complexity of the HPC storage stack rapidly grows, domain scientists face increasing challenges in effectively utilizing HPC storage systems to achieve their desired I/O performance. To identify and address I/O issues, scientists largely rely on I/O experts to analyze their I/O traces and provide insights into potential problems. However, with a limited number of I/O experts and the growing demand for dataintensive applications, inaccessibility has become a major bottleneck, hindering scientists from maximizing their productivity. The recent rapid progress in large language models (LLMs) opens the door to creating an automated tool that democratizes trustworthy I/O performance diagnosis capabilities to domain scientists. However, LLMs face significant challenges in this task, such as the inability to handle long context windows, a lack of accurate domain knowledge about HPC I/O, and the generation of hallucinations during complex interactions. In this work, we propose IOAgent as a systematic effort to address these challenges. IOAgent integrates various new designs, including a module-based pre-processor, a RAG-based domain knowledge integrator, and a tree-based merger to accurately diagnose I/O issues from a given Darshan trace file. Similar to an I/O expert, IOAgent provides detailed justifications and references for its diagnoses and offers an interactive interface for scientists to continue asking questions about the diagnosis. To evaluate IOAgent, we collected a diverse set of labeled job traces and released the first open diagnosis test suite, TraceBench. Based on this test suite, extensive evaluations were conducted, demonstrating that IOAgent matches or outperforms state-of-the-art I/O diagnosis tools with accurate and useful diagnosis results. We also show that IOAgent is not tied to specific LLMs, performing similarly well with both proprietary and open-source LLMs. We believe IOAgent has the potential to become a powerful tool for scientists navigating complex HPC I/O subsystems in the future. Chris Egersdoerfer, Arnav Sareen, Jean Luca Bez, Surendra Byna, Dongkuan Xu, Dong Dai 0001 |
IPDPS | 6 |
| 2025 | Be Aware of Metadata Corruption in Parallel File System: It can be Silent and CatastrophicabstractHigh-Performance Computing (HPC) systems rely on parallel file systems (PFS) like Lustre to reliably manage large-scale data. Unfortunately, such data stored in PFS are subjected to corruption due to hardware failures, power outages, on-disk corruption, in-memory corruption, and software bugs. Previous work has shown the impacts of general data corruptions on scientific applications, but barely investigated the impacts of corruptions on metadata, which is a special type of data for many critical PFS functionalities. Since the metadata are accessed and updated frequently, the metadata corruption is non-negligible, and the impact on both HPC systems and applications is often complicated. In this study, we systematically studied the effects of PFS metadata corruption on representative scientific applications and workflows through fault injections. We observed various abnormal behaviors of applications against metadata corruptions, such as failing to finish execution or successfully finishing execution but generating wrong outputs silently. We also showed that neither the defensive programming practice from the application developers nor existing metadata checking mechanism deployed in PFS can resolve the metadata corruption issues or guild applications to react correctly. To address this issue, we further propose a metadata checksum mechanism as a system-level mitigation strategy to detect and handle metadata corruptions at runtime. We implement a prototype of the checksum mechanism on Lustre file system using FUSE interface and evaluate its functionality and performance overhead. Our results show the proposed checksum mechanism can effectively detect metadata corruptions during data accesses with extremely low overhead and stop the applications from continuous execution and generating wrong results. Saisha Kamat, Mai Zheng, Bo Fang 0002, Dong Dai 0001 |
IPDPS | 4 |
| 2025 | AdapTBF: Decentralized Bandwidth Control via Adaptive Token Borrowing for HPC StorageabstractModern high-performance computing (HPC) applications run exclusively on computational resources but share global storage systems. This design can lead to issues when applications use disproportional amount of storage resources compared to their allocated computing resources. An application running on a single compute node might consume excessive I/O bandwidth from a storage server, for example, by issuing numerous small, random writes. In doing so, it can hinder larger jobs that also write to the same storage server and are allocated many compute nodes, resulting in significant resource waste. A straightforward solution to prevent such an issue is to limit each application's I/O bandwidth on storage servers in proportion to its allocated computation resources. This approach has been implemented in existing parallel file systems using the Token Bucket Filter (TBF) mechanism. However, applying such methods in practice often results in lower overall I/O efficiency. HPC applications are known for generating short, bursty I/O requests. Therefore, strictly limiting I/O bandwidth proportionally either leads to wasted storage server bandwidth when applications are not performing I/O, or it prevents applications from temporarily utilizing higher bandwidth during bursty I/O phases. We argue that the goal of I/O control should be to maximize both the I/O performance of each application and the overall efficiency of the storage server, while ensuring fairness among jobs, such as preventing smaller jobs from blocking largescale ones. In this paper, we propose AdapTBF to achieve this objective. Building on the well-established TBF mechanism in modern parallel file systems (e.g., Lustre), AdapTBF introduces a decentralized bandwidth control approach using an adaptive borrowing and lending mechanism. We detail the algorithm, implement AdapTBF on the Lustre file system, and evaluate it using synthetic workloads modeled after real-world scenarios to demonstrate its effectiveness. Experimental results show that AdapTBF effectively manages I/O bandwidth for applications on storage servers while maintaining high overall storage utilization, even under extreme conditions. Md. Hasanur Rashid, Dong Dai 0001 |
IPDPS | 2 |
| 2025 | STELLAR: Storage Tuning Engine Leveraging LLM Autonomous Reasoning for High Performance Parallel File SystemsabstractI/O performance is crucial to efficiency in data-intensive scientific computing; but tuning large-scale storage systems is complex, costly, and notoriously manpower-intensive, making it inaccessible for most domain scientists. To address this problem, we propose STELLAR, an autonomous tuner for high-performance parallel file systems. Our evaluations show that STELLAR almost always selects near-optimal configurations for the parallel file systems within the first five attempts, even for previously unseen applications. STELLAR’s human-like efficiency is fundamentally different from existing autotuning methods, which often require hundreds of thousands of iterations to converge. STELLAR achieves this through autonomous end-to-end agentic tuning. Powered by large language models (LLMs), STELLAR is capable of (1) accurately extracting tunable parameters from software manuals, (2) analyzing I/O trace logs generated by applications, (3) selecting initial tuning strategies, (4) rerunning applications on real systems and collecting I/O performance feedback, (5) adjusting tuning strategies and repeating the tuning cycle, and (6) reflecting on and summarizing tuning experiences into reusable knowledge for future optimizations. STELLAR integrates retrieval-augmented generation (RAG), external tool execution, LLM-based reasoning, and a multiagent design to stabilize reasoning and combat hallucinations. We evaluate how each of these components impacts optimization outcomes, thus providing insight into the design of similar systems for other optimization problems. STELLAR’s architecture and empirical validation open new avenues for tackling complex system optimization challenges, especially those characterized by vast search spaces and high exploration costs. Its highly efficient autonomous tuning will broaden access to I/O performance optimizations for domain scientists with minimal additional resource investment. Chris Egersdoerfer, Philip H. Carns, Shane Snyder, Robert B. Ross, Dong Dai 0001 |
SC | 5 |
| 2025 | Improving SpGEMM Performance Through Matrix-Reordering and Cluster-wise ComputationabstractSparse matrix-sparse matrix multiplication (SpGEMM) is a key kernel in many scientific applications and graph workloads. Unfortunately, SpGEMM is bottlenecked by data movement due to its irregular memory access patterns. Significant work has been devoted to developing row reordering schemes towards improving locality in sparse operations, but prior studies mostly focus on the case of sparse-matrix vector multiplication (SpMV). Abdullah Al Raqibul Islam, Helen Xu 0001, Dong Dai 0001, Aydin Buluç |
SC | 3 |
| 2025 | Hardware Accelerated Vision Transformer via Heterogeneous Architecture Design and Adaptive Dataflow MappingabstractVision transformer (ViT) models have demonstrated remarkable advantages in visual tasks. However, the ViT model contains various types of operators, and its sophisticated model structure imposes substantial computational complexity and storage burden. Existing hardware solutions still fail to fully unleash the ViT acceleration potential due to the mismatch between operators and hardware architectures, suffering from inefficient dataflow mapping. This work proposes HDViT, a full-fledged heterogeneous hardware accelerator on FPGA, to enhance the ViT acceleration by comprehensively analyzing and addressing the challenges of heterogeneous architecture design. Specifically, HDViT first develops a heterogeneous architecture design that is composed of multiple processing engines (PEs) to accelerate various operators in the ViT model. Then, HDViT devises a hybrid-oriented dataflow mapping strategy to reduce data transmission granularity and alleviate storage resource pressure. Lastly, to achieve the latency balancing among multiple PEs, we formulate the HDViT architecture and implement an automated exploration process to identify optimized parallelism parameters that satisfy computation and storage demands while enhancing the heterogeneous architectural performance. Experimental results indicate that HDViT achieves significant performance speedups of 2.16$\times$and 3.51$\times$compared to previous heterogeneous and unified accelerators, respectively. HDViT also achieves a maximum of 98.46% hardware utilization. Yingxue Gao, Lei Gong 0003, Chao Wang 0003, Dong Dai 0001, Yang Yang 0080, Xianglan Chen, Xi Li 0003, Xuehai Zhou |
IEEE Trans. Computers | 5 |
| 2024 | QualityNet: Error-bounded Lossy Compression Quality Prediction via Deep SurrogateabstractAs scientific simulations generate increasingly large datasets, efficient data compression becomes essential to mitigate storage and transmission challenges. Error-bounded Lossy compression algorithms reduce data volume at the expense of data fidelity. However, the evaluation of compressed data quality can be resource-intensive and time-consuming. This paper proposes a surrogate-based framework to predict key compression quality assessment metrics, allowing users to evaluate the quality of compressed data quickly without extensive trial and error. We evaluate our proposed framework on real-world scientific application datasets. Our surrogate-based framework shows superior prediction accuracy, with the RMSE generally lower than 5%. On the other hand, our framework demonstrates robust generalization for prediction across data fields within the same application. Our results show that our approach significantly reduces the error bound selection time and accelerates the downstream evaluation process for decompressed data, making it easier for domain scientists to select optimal error bounds of lossy compression while maintaining data quality that meets their needs. This work bridges the gap between compression ratio optimization and data fidelity, offering a scalable solution for scientific applications that rely on large-scale lossy compression. Khondoker Mirazul Mumenin, Dong Dai 0001, Jinzhen Wang, Sheng Di |
IEEE Big Data | 2 |
| 2024 | ION: Navigating the HPC I/O Optimization Journey using Large Language ModelsabstractEffectively leveraging the complex software and hardware I/O stacks of HPC systems to deliver needed I/O performance has been a challenging task for domain scientists. To identify and address I/O issues in their applications, scientists largely rely on I/O experts to analyze the recorded I/O traces of their applications and provide insights into the potential issues. However, due to the limited number of I/O experts and the growing demand for data-intensive applications across the wide spectrum of sciences, inaccessibility has become a major bottleneck hindering scientists from maximizing their productivity. Inspired by the recent rapid progress of large language models (LLMs), in this work we propose IO Navigator (ION), an LLM-based framework that takes a recorded I/O trace of an application as input and leverages the in-context learning, chain-of-thought, and code generation capabilities of LLMs to comprehensively analyze the I/O trace and provide diagnosis of potential I/O issues. Similar to an I/O expert, ION provides detailed justifications for the diagnosis and an interactive interface for scientists to ask detailed questions about the diagnosis. We illustrate ION's applicability by assessing it on a set of controlled I/O traces generated with different I/O issues. We also demonstrate that ION can match state-of-the-art I/O optimization tools and provide more insightful and adaptive diagnoses for real applications. We believe ION, with its full capabilities, has the potential to become a powerful tool for scientists to navigate through complex I/O subsystems in the future. Chris Egersdoerfer, Arnav Sareen, Jean Luca Bez, Surendra Byna, Dong Dai 0001 |
HotStorage | 5 |
| 2024 | Cross-System Analysis of Job Characterization and Scheduling in Large-Scale Computing ClustersabstractAmid the growing prevalence of artificial intelligence (AI) and deep learning (DL) across industries and science disciplines, high-performance computing (HPC) clusters are increasingly used for DL tasks, in addition to their traditional role in numerical simulations. This shared use of HPC systems for both DL tasks and numerical codes is altering the characteristics of their workloads, leaving many previously observed and well-accepted workload characteristics unchecked, potentially outdated and imprecise, for these new mixed workloads. Thus, to understand these changes and their implications for job scheduling, we conduct a cross-system analysis of job characterization and its scheduling across a range of representative clusters, including two classic HPC clusters (Mira, Theta), two classic DL clusters (Philly, Helios), and a hybrid cluster (Blue Waters).Our cross-system analyses focus on three key aspects: 1) job geometries (job size, run time, arrival interval) and their impacts on scheduling results, 2) job failure patterns and their correlations to job geometries, 3) per-user behaviors and their indication on job scheduling. Through these comparisons, we confirm notable disparity and similarity among different systems (summarized as 8 takeaways), which would help design more efficient job schedulers for the future HPC systems. We further introduce two use case studies (job runtime prediction and adaptive relaxed backfilling) that leverage these new observations and show the improved job scheduling results. In summary, we expect our observations, insights, and systematic analysis approach can be useful for the community in building efficient HPC scheduling for the upcoming hybrid workloads. Di Zhang 0015, Monish Soundar Raj, Sheng Di, Dong Dai 0001 |
IPDPS | 5 |
| 2024 | An Empirical Study of Machine Learning-Based Synthetic Job Trace Generation Methods
Monish Soundar Raj, Thomas MacDougall, Di Zhang 0015, Dong Dai 0001 |
JSSPP | 4 |
| 2024 | PROV-IO$^+$+: A Cross-Platform Provenance Framework for Scientific Data on HPC SystemsabstractData provenance, or data lineage, describes the life cycle of data. In scientific workflows on HPC systems, scientists often seek diverse provenance (e.g., origins of data products, usage patterns of datasets). Unfortunately, existing provenance solutions cannot address the challenges due to their incompatible provenance models and/or system implementations. In this paper, we analyze four representative scientific workflows in collaboration with the domain scientists to identify concrete provenance needs. Based on the first-hand analysis, we propose a provenance framework called PROV-IO$^+$, which includes an I/O-centric provenance model for describing scientific data and the associated I/O operations and environments precisely. Moreover, we build a prototype of PROV-IO$^+$to enable end-to-end provenance support on real HPC systems with little manual effort. The PROV-IO$^+$framework can support both containerized and non-containerized workflows on different HPC platforms with flexibility in selecting various classes of provenance. Our experiments with realistic workflows show that PROV-IO$^+$can address the provenance needs of the domain scientists effectively with reasonable performance (e.g., less than 3.5% tracking overhead for most experiments). Moreover, PROV-IO$^+$outperforms a state-of-the-art system (i.e., ProvLake) in our experiments. Runzhou Han, Mai Zheng, Surendra Byna, Houjun Tang, Bin Dong 0002, Dong Dai 0001, Yong Chen 0001, Dongkyun Kim, Joseph Hassoun, David Thorsley |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2023 | Early Exploration of Using ChatGPT for Log-based Anomaly Detection on Parallel File Systems LogsabstractLog-based anomaly detection has been extensively studied to help detect complex runtime anomalies in production systems. However, existing techniques exhibit several common issues. First, they rely heavily on expert-labeled logs to discern anomalous behavior patterns. But labelling enough log data manually to effectively train deep neural networks may take too long. Second, they rely on numeric model prediction based on numeric vector input which causes model decisions to be largely non-interpretable by humans which further rules out targeted error correction. Chris Egersdoerfer, Di Zhang 0015, Dong Dai 0001 |
HPDC | 3 |
| 2023 | FaultyRank: A Graph-based Parallel File System CheckerabstractSimilar to local file system checkers such as e2fsck for Ext4, a parallel file system (PFS) checker ensures the file system's correctness. The basic idea of file system checkers is straightforward: important metadata are stored redundantly in separate places for cross-checking; inconsistent metadata will be repaired or overwritten by its ‘more correct' counterpart, which is defined by the developers. Unfortunately, implementing the idea for PFSes is non-trivial due to the system complexity. Although many popular parallel file systems already contain dedicated checkers (e.g., LFSCK for Lustre, BeeGFS-FSCK for BeeGFS, mmfsck for GPFS), the existing checkers often cannot detect or repair inconsistencies accurately due to one fundamental limitation: they rely on a fixed set of consistency rules predefined by developers, which cannot cover the various failure scenarios that may occur in practice.In this study, we propose a new graph-based method to build PFS checkers. Specifically, we model important PFS metadata into graphs, then generalize the logic of cross-checking and repairing into graph analytic tasks. We design a new graph algorithm, FaultyRank, to quantitatively calculate the correctness of each metadata object. By leveraging the calculated correctness, we are able to recommend the most promising repairs to users. Based on the idea, we implement a prototype of FaultyRank on Lustre, one of the most widely used parallel file systems, and compare it with Lustre's default file system checker LFSCK. Our experiments show that FaultyRank can achieve the same checking and repairing logic as LFSCK. Moreover, it is capable of detecting and repairing complicated PFS consistency issues that LFSCK can not handle. We also show the performance advantage of FaultyRank compared with LFSCK. Through this study, we believe FaultyRank opens a new opportunity for building PFS checkers effectively and efficiently. Saisha Kamat, Abdullah Al Raqibul Islam, Mai Zheng, Dong Dai 0001 |
IPDPS | 4 |
| 2023 | Drill: Log-based Anomaly Detection for Large-scale Storage Systems Using Source Code AnalysisabstractLarge-scale storage systems, a critical part of modern computing systems, are subject to various runtime bugs, failures, and anomalies in production. Identifying their anomalies at runtime is thus critical for users and administrators. Since runtime logs record the important status of the systems, log-based anomaly detection has been studied extensively for timely identifying system malfunctions. However, existing log-based anomaly detection solutions share common limitations in representing log entries accurately and robustly, hence can not effectively handle log entries that were not seen in the historical logs, which is a common real-world scenario due to logs' inherent rarity and the continuous evolution of the systems. To address the issues of existing methods, we propose Drill, a new log pre-processing method to generate high-quality vector representation of runtime logs by leveraging both storage system-specific sentiment-classifying language models and log contexts built from the source code. Through extensive evaluations of two representative distributed storage systems (Apache HDFS and Lustre), we show that Drill can achieve up to 41% improvement when compared with state-of-the-art anomaly detection solutions, showing it is a promising solution for general anomaly detection. Di Zhang 0015, Chris Egersdoerfer, Tabassum Mahmud, Mai Zheng, Dong Dai 0001 |
IPDPS | 5 |
| 2023 | DGAP: Efficient Dynamic Graph Analysis on Persistent MemoryabstractDynamic graphs, featuring continuously updated vertices and edges, have grown in importance for numerous real-world applications. To accommodate this, graph frameworks, particularly their internal data structures, must support both persistent graph updates and rapid graph analysis simultaneously, leading to complex designs to orchestrate 'fast but volatile' and 'persistent but slow' storage devices. Emerging persistent memory technologies, such as Optane DCPMM, offer a promising alternative to simplify the designs by providing data persistence, low latency, and high IOPS together. In light of this, we propose DGAP, a framework for efficient dynamic graph analysis on persistent memory. Unlike traditional dynamic graph frameworks, which combine multiple graph data structures (e.g., edge list or adjacency list) to achieve the required performance, DGAP utilizes a single mutable Compressed Sparse Row (CSR) graph structure with new designs for persistent memory to construct the framework. Specifically, DGAP introduces a per-section edge log to reduce write amplification on persistent memory; a per-thread undo log to enable high-performance, crash-consistent rebalancing operations; and a data placement schema to minimize in-place updates on persistent memory. Our extensive evaluation results demonstrate that DGAP can achieve up to 3.2× better graph update performance and up to 3.77× better graph analysis performance compared to state-of-the-art dynamic graph frameworks for persistent memory, such as XPGraph, LLAMA, and GraphOne. Abdullah Al Raqibul Islam, Dong Dai 0001 |
SC | 2 |
| 2023 | Dynamic Resource Provisioning for Iterative Workloads on Apache SparkabstractApache Spark as a popular in-memory data analytic framework has been employed by various applications—such as machine learning, graph computation, and scientific computing, which benefit from the long-running process (e.g., executor) programming model to avoid system I/O overhead. However, existing resource allocation strategies mainly rely on the peak demand, which are normally specified by users. Since the resource usages of long-running applications like iterative computation vary significantly over time, we find that peak demand based resource allocation policies lead to low cloud utilization in production environments. In this article, we present a utilization aware resource provisioning approach for iterative workloads on Apache Spark (i.e.,${iSpark}$). It can identify the causes of resource underutilization due to an inflexible resource policy, and elastically adjusts the allocated executors over time according to the real-time resource usage. In general, iterative applications require more computation resources at the beginning stage and their demands for resources diminish as more iterations are completed. iSpark aims to timely scale up or scale down the number of executors in order to fully utilize the allocated resources while taking the dominant factor into consideration. It further preempts the underutilized executors and preserves the cached intermediate data to ensure the data consistency. Testbed evaluations show that iSpark averagely improves the resource utilization of individual executors by 35.2% compared to vanilla Spark. At the same time, it increases the cluster utilization from 32.1% to 51.3% and effectively reduces the overall job completion time by 20.8% for a set of representative iterative applications. Furthermore, we have extended iSpark to multi-tenancy cloud environments. Specifically, iSpark characterizes a virtual node based on two real-time measured performance statistics: I/O rate and CPU steal time. Thus, we extend the two-dimensional resource constraints (i.e., CPU and MeM) in iSpark to three-dimensional resource constraints (i.e., CPU, MeM and I/O) in the cloud environment. We consider two representative interference scenarios in the cloud: stable interference and dynamic interference. Experimental results on virtual clusters with varying interferences show that iSpark with cloud extension improves the average job completion time by 68% compared to default Spark resource allocation policies. Dazhao Cheng, Yu Wang 0003, Dong Dai 0001 |
IEEE Trans. Cloud Comput. | 3 |
| 2022 | VCSR: Mutable CSR Graph Format Using Vertex-Centric Packed Memory ArrayabstractThe compressed sparse row (CSR) is a widely used graph storage format due to its compact memory layout and high performance on graph analytic tasks. However, the compact design also limits itself from supporting many graph applications that operate on dynamic or temporal graphs as updates on these graph will need to rebuild the entire CSR structure, leading to high costs. Extending CSR to support efficient graph mutations without losing its high performance on graph analysis then becomes critical to these applications. Existing mutable CSR extensions leverage packed memory array (PMA) to store edge list to enable graph mutations. But, such a naive way has fundamental limitations in handling imbalanced graphs, which many real-world graphs belong to. To address such issues, we propose VCSR, a new mutable CSR storage format that leverages the packed memory array (PMA) via a new vertex-centric strategy to efficiently support temporal graphs. Our evaluation results show that compared with the state-of-the-art mutable CSR extensions, VCSR can achieve 1.41x-3.81x better performance in graph insertions and 1.22x-2.05x better performance in running typical graph analytic algorithms. In addition, VCSR can achieve similar performance as the original immutable CSR in running graph analytic tasks, making it a promising storage format for temporal graphs. Abdullah Al Raqibul Islam, Dong Dai 0001, Dazhao Cheng |
CCGRID | 2 |
| 2022 | SchedInspector: A Batch Job Scheduling Inspector Using Reinforcement LearningabstractImproving the performance of job executions is an important goal of HPC batch job schedulers, such as minimizing job waiting time, slowdown, or completion time. Such a goal is often accomplished using carefully designed heuristics based on job features, such as job size and job duration. However, these heuristics overlook important runtime factors (e.g., cluster availability and waiting job patterns), which may vary across time and make a previously sound scheduling decision not hold any longer. In this study, we propose a new approach to incorporate runtime factors into batch job scheduling for better job execution performance. The key idea is to add a scheduling inspector on top of the base job scheduler to scrutinize its scheduling decisions. The inspector will take the runtime factors into consideration and accordingly determine the fitness of the scheduled job. It then either accepts the scheduled job or rejects it and asks the base schedulers to try again later. We realize such an inspector, namely SchedInspector, by leveraging the intelligence of reinforcement learning. Through extensive experiments, we show SchedInspector can intelligently integrate the runtime factors into various batch job scheduling policies, including the state-of-the-art one, to gain better job execution performance, such as smaller average bounded job slowdown (up to 69% better) or average job waiting time (up to 52% better), across various real-world workloads. We also show that although rejecting scheduling decisions may leave the resources idle hence affect the system utilization, SchedInspector is able to achieve the job execution performance improvement with marginal impact on the system utilization (typically less than 1%). We consider one key advantage of SchedInspector is it automatically learns to work with and improve existing job scheduling policies without changing them, which makes it promising to serve as a generic enhancer for various batch job scheduling policies. Di Zhang 0015, Dong Dai 0001 |
HPDC | 2 |
| 2022 | A performance study of optane persistent memory: from storage data structures' perspective
Abdullah Al Raqibul Islam, Christopher York, Dong Dai 0001 |
CCF Trans. High Perform. Comput. | 3 |
| 2022 | A Study of Failure Recovery and Logging of High-Performance Parallel File SystemsabstractLarge-scale parallel file systems (PFSs) play an essential role in high-performance computing (HPC). However, despite their importance, their reliability is much less studied or understood compared with that of local storage systems or cloud storage systems. Recent failure incidents at real HPC centers have exposed the latent defects in PFS clusters as well as the urgent need for a systematic analysis. To address the challenge, we perform a study of the failure recovery and logging mechanisms of PFSs in this article. First, to trigger the failure recovery and logging operations of the target PFS, we introduce a black-box fault injection tool called PFault , which is transparent to PFSs and easy to deploy in practice. PFault emulates the failure state of individual storage nodes in the PFS based on a set of pre-defined fault models and enables examining the PFS behavior under fault systematically. Next, we apply PFault to study two widely used PFSs: Lustre and BeeGFS. Our analysis reveals the unique failure recovery and logging patterns of the target PFSs and identifies multiple cases where the PFSs are imperfect in terms of failure handling. For example, Lustre includes a recovery component called LFSCK to detect and fix PFS-level inconsistencies, but we find that LFSCK itself may hang or trigger kernel panics when scanning a corrupted Lustre. Even after the recovery attempt of LFSCK, the subsequent workloads applied to Lustre may still behave abnormally (e.g., hang or report I/O errors). Similar issues have also been observed in BeeGFS and its recovery component BeeGFS-FSCK. We analyze the root causes of the abnormal symptoms observed in depth, which has led to a new patch set to be merged into the coming Lustre release. In addition, we characterize the extensive logs generated in the experiments in detail and identify the unique patterns and limitations of PFSs in terms of failure logging. We hope this study and the resulting tool and dataset can facilitate follow-up research in the communities and help improve PFSs for reliable high-performance computing. Runzhou Han, Om Rameshwar Gatla, Mai Zheng, Jinrui Cao, Di Zhang 0015, Dong Dai 0001, Yong Chen 0001, Jonathan E. Cook 0001 |
ACM Trans. Storage | 6 |
| 2021 | SentiLog: Anomaly Detecting on Parallel File Systems via Log-based Sentiment AnalysisabstractAs core components of High-performance computing (HPC) platforms, parallel file systems (PFSes) grow quickly in scale and complexity, hence are subject to various failures and anomalies. Identifying their anomalies in runtime is critically helpful for HPC operators and administrators. Analyzing the runtime logs to detect the anomalies of large-scale systems has been proven effective in many recent studies. However, applying them to parallel file systems logs faces significant challenges due to the large volume and irregularity of PFSes logs. This study proposes SentiLog, a new approach to analyzing PFSes system logs for detecting anomalies. Unlike existing solutions, SentiLog works by training a general sentimental, natural language model based on the logging-relevant source code collected from a set of PFSes. In this way, SentiLog learns information embedded by developers from the source code. Our preliminary results show SentiLog is able to accurately predict anomalies and performs better than state-of-the-art log analysis solutions on two representative PFSes (Lustre and BeeGFS). This preliminary study shows sentiment analysis could be a promising method to analyze complex and irregular system logs. Di Zhang 0015, Dong Dai 0001, Runzhou Han, Mai Zheng |
HotStorage | 2 |
| 2021 | I/O characteristic discovery for storage system optimizations
Yong Chen 0001, Dong Dai 0001, Weiping Wang 0005 |
J. Parallel Distributed Comput. | 3 |
| 2021 | Trigger-Based Incremental Data Processing with Unified Sync and Async ModelabstractIn recent years, more and more applications in the cloud have needs to process large-scale on-line datasets, which evolve over time as new entries are added and existing entries are modified. Several programming frameworks, such as Percolator and Oolong, are proposed for such incremental data processing and can achieve efficient processing with an event-driven abstraction. However, these frameworks are inherently asynchronous, leaving the heavy burden of managing synchronization to applications' developers, which further significantly restricts their usabilities. In this study, we propose a trigger-based incremental computing framework in the cloud, called Domino, with both synchronous and asynchronous mechanisms to coordinate parallel triggers. With this new framework, both synchronous and asynchronous applications can be seamlessly developed. Use cases and extensive evaluation results confirm that it can deliver sufficient performance, and also is easy to use for incremental applications in large-scale distributed computing. Dong Dai 0001, Yong Chen 0001, Dries Kimpe, Robert B. Ross |
IEEE Trans. Cloud Comput. | 1 |
| 2020 | Understand the overheads of storage data structures on persistent memoryabstractThe byte-addressable persistent memory (PMEM) devices have opened new opportunities for building high-performance storage systems. With both DRAM and PMEM in the system, it is important to choose the correct storage data structures on each of them to achieve the best overall performance and the needed data persistence. However, this is non-trivial. One reason is the limited understanding of the actual performance characteristic of different data structures on PMEM. In this study, we develop storeds-bench to help developers better understand the overhead they will encounter when using a certain data structure on PMEM. Specifically, storeds-bench is designed as a benchmark suite that leverages YCSB and has various commonly used storage data structures implemented using PMDK (persistent memory develop kit) under different persistent and consistency requirements. Abdullah Al Raqibul Islam, Dong Dai 0001 |
PPoPP | 2 |
| 2020 | RLScheduler: an automated HPC batch job scheduler using reinforcement learningabstractToday's high-performance computing (HPC) platforms are still dominated by batch jobs. Accordingly, effective batch job scheduling is crucial to obtain high system efficiency. Existing HPC batch job schedulers typically leverage heuristic priority functions to prioritize and schedule jobs. But, once configured and deployed by the experts, such priority functions can hardly adapt to the changes of job loads, optimization goals, or system settings, potentially leading to degraded system efficiency when changes occur. To address this fundamental issue, we present RLScheduler, an automated HPC batch job scheduler built on reinforcement learning. RLScheduler relies on minimal manual interventions or expert knowledge, but can learn high-quality scheduling policies via its own continuous `trial and error'. We introduce a new kernel-based neural network structure and trajectory filtering mechanism in RLScheduler to improve and stabilize the learning process. Through extensive evaluations, we confirm that RLScheduler can learn high-quality scheduling policies towards various workloads and various optimization goals with relatively low computation cost. Moreover, we show that the learned models perform stably even when applied to unseen workloads, making them practical for production use. Di Zhang 0015, Dong Dai 0001, Youbiao He, Forrest Sheng Bao |
SC | 2 |
| 2020 | PRS: A Pattern-Directed Replication Scheme for Heterogeneous Object-Based StorageabstractData replication is a key technique to achieve high data availability, reliability, and optimized performance in distributed storage systems. In recent years, with emerged new storage devices, heterogeneous object-based storage systems, such as a storage system with a mix of hard disk drives, solid state drives, and other non-volatile memory devices have become increasingly attractive since they combine the merits of different storage devices to deliver better promises. However, existing data replication schemes do not well consider distinct characteristics of heterogeneous storage devices yet, which could lead to suboptimal performance. This article introduces a new data replication scheme called Pattern-directed Replication Scheme (PRS) to achieve efficient data replication for heterogeneous storage systems. Different from traditional schemes, the PRS selectively replicates data objects and distributes replicas to various storage devices based on their characteristics. It aggregates objects that have I/O correlation into object groups by calculating object distance and makes replication for grouped objects according to application's data access pattern identified. In addition, the PRS uses a pseudo random algorithm to optimize replica placement by considering the storage device performance and capacity features. We have evaluated the pattern-directed replication scheme with extensive tests in Sheepdog, a typical object-based storage system. The experimental results confirm that it is a highly efficient replication scheme for heterogeneous storage systems. For instance, the read performance was improved by 105 percent to nearly 10x compared with existing replication schemes. Yong Chen 0001, Wei Xie 0017, Dong Dai 0001, Shuibing He, Weiping Wang 0005 |
IEEE Trans. Computers | 4 |
| 2019 | A Performance Study of Lustre File System Checker: Bottlenecks and PotentialsabstractLustre, as one of the most popular parallel file systems in high-performance computing (HPC), provides POSIX interface and maintains a large set of POSIX-related metadata, which could be corrupted due to hardware failures, software bugs, configuration errors, etc. The Lustre file system checker (LFSCK) is the remedy tool to detect metadata inconsistencies and to restore a corrupted Lustre to a valid state, hence is critical for reliable HPC. Unfortunately, in practice, LFSCK runs slow in large deployment, making system administrators reluctant to use it as a routine maintenance tool. Consequently, cascading errors may lead to unrecoverable failures, resulting in significant downtime or even data loss. Given the fact that HPC is rapidly marching to Exascale and much larger Lustre file systems are being deployed, it is critical to understand the performance of LFSCK. In this paper, we study the performance of LFSCK to identify its bottlenecks and analyze its performance potentials. Specifically, we design an aging method based on real-world HPC workloads to age Lustre to representative states, and then systematically evaluate and analyze how LFSCK runs on such an aged Lustre via monitoring the utilization of various resources. From our experiments, we find out that the design and implementation of LFSCK is sub-optimal. It consists of scalability bottleneck on the metadata server (MDS), relatively high fan-out ratio in network utilization, and unnecessary blocking among internal components. Based on these observations, we discussed potential optimization and present some preliminary results. Dong Dai 0001, Om Rameshwar Gatla, Mai Zheng |
MSST | 1 |
| 2019 | Vectorizing disks blocks for efficient storage system via deep learning
Dong Dai 0001, Forrest Sheng Bao, Xuanhua Shi, Yong Chen 0001 |
Parallel Comput. | 1 |
| 2019 | Client-side straggler-aware I/O scheduler for object-based parallel file systems
Neda Tavakoli, Dong Dai 0001, Yong Chen 0001 |
Parallel Comput. | 2 |
| 2019 | Managing Rich Metadata in High-Performance Computing Systems Using a Graph ModelabstractHigh-performance computing (HPC) systems generate huge amounts of metadata about different entities such as jobs, users, and files. Existing systems can efficiently record and manage part of these metadata, mainly the POSIX metadata of data files (e.g., file size, name, and permissions mode). But another important set of metadata, referred to as “rich” metadata in this study, which record not only wider range of entities (e.g., running processes and jobs) but also more complex relationships between them, are mostly missing in current HPC systems. Yet such rich metadata are critical for supporting many advanced data management functions such as identifying data sources and parameters behind a given result; auditing data usage; or understanding details about how inputs are transformed into outputs. To uniformly and efficiently manage the rich metadata generated in HPC systems, We propose to utilize a graph model in this study. We identify the key challenges of implementing such a graph-based HPC rich metadata management system and present GraphMeta, a graph-based rich metadata management system designed and optimized for HPC platforms, to tackle these challenges. Extensive evaluations on both synthetic and real HPC metadata workloads show its advantages in both performance and scalability compared with existing solutions. Dong Dai 0001, Yong Chen 0001, Philip H. Carns, John Jenkins, Wei Zhang 0097, Robert B. Ross |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2018 | I/O Characteristics Discovery in Cloud Storage SystemsabstractThe data growth from many applications in clouds poses significant challenges to cloud storage systems. To deliver the best storage and I/O performance possible, it is often required to understand and leverage the I/O characteristics based on data accesses. A number of research studies have been carried out on this topic. However, most of them either utilize a limited number of data-access attributes, restricting the general applicability of the method for different applications, or heavily rely on the domain knowledge or expertise about applications' I/O behaviors to select the best representative features, introducing bias for certain workloads. To overcome these limitations, in this study, we present a new I/O characteristic discovery methodology. This method enables capturing data-access features as many as possible to eliminate human bias. It utilizes a machine-learning based strategy to derive the most important set of features automatically, and groups data objects with a clustering algorithm (DBSCAN) to reveal I/O characteristics discovered. These I/O characteristics revealed can direct I/O performance optimizations in numerous scenarios, such as in data prefeteching and data reorganization optimizations in cloud storage systems. Dong Dai 0001, Yong Chen 0001 |
IEEE CLOUD | 2 |
| 2018 | AKIN: A Streaming Graph Partitioning Algorithm for Distributed Graph Storage SystemsabstractMany graph-related applications face the challenge of managing excessive and ever-growing graph data in a distributed environment. Therefore, it is necessary to consider a graph partitioning algorithm to distribute graph data onto multiple machines as the data comes in. Balancing data distribution and minimizing edge-cut ratio are two basic pursuits of the graph partitioning problem. While achieving balanced partitions for streaming graphs is easy, existing graph partitioning algorithms either fail to work on streaming workloads, or leave edge-cut ratio to be further improved. Our research aims to provide a better solution that fits the need of streaming graph partitioning in a distributed system, which further reduces the edge-cut ratio while maintaining rough balance among all partitions. We exploit the similarity measure on the degree of vertices to gather structuralrelated vertices in the same partition as much as possible, this reduces the edge-cut ratio even further as compared to the state-of-the-art streaming graph partitioning algorithm - FENNEL. Our evaluation shows that our streaming graph partitioning algorithm is able to achieve better partitioning quality in terms of edge-cut ratio (up to 20% reduction as compared to FENNEL) while maintaining decent balance between all partitions, and such improvement applies to various real-life graphs. Wei Zhang 0097, Yong Chen 0001, Dong Dai 0001 |
CCGrid | 3 |
| 2018 | PFault: A General Framework for Analyzing the Reliability of High-Performance Parallel File SystemsabstractHigh-performance parallel file systems (PFSes) are of prime importance today. However, despite the importance, their reliability is much less studied compared with that of local storage systems, largely due to the lack of an effective analysis methodology. Jinrui Cao, Om Rameshwar Gatla, Mai Zheng, Dong Dai 0001, Vidya Eswarappa, Yan Mu, Yong Chen 0001 |
ICS | 4 |
| 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 | 6 |
| 2017 | Lightweight Provenance Service for High-Performance ComputingabstractProvenance describes detailed information about the history of a piece of data, containing the relationships among elements such as users, processes, jobs, and workflows that contribute to the existence of data. Provenance is key to supporting many data management functionalities that are increasingly important in operations such as identifying data sources, parameters, or assumptions behind a given result; auditing data usage; or understanding details about how inputs are transformed into outputs. Despite its importance, however, provenance support is largely underdeveloped in highly parallel architectures and systems. One major challenge is the demanding requirements of providing provenance service in situ. The need to remain lightweight and to be always on often conflicts with the need to be transparent and offer an accurate catalog of details regarding the applications and systems. To tackle this challenge, we introduce a lightweight provenance service, called LPS, for high-performance computing (HPC) systems. LPS leverages a kernel instrument mechanism to achieve transparency and introduces representative execution and flexible granularity to capture comprehensive provenance with controllable overhead. Extensive evaluations and use cases have confirmed its efficiency and usability. We believe that LPS can be integrated into current and future HPC systems to support a variety of data management needs. Dong Dai 0001, Yong Chen 0001, Philip H. Carns, John Jenkins, Robert B. Ross |
PACT | 1 |
| 2017 | Pattern-Directed Replication Scheme for Heterogeneous Object-based StorageabstractData replication is a key technique to achieve data availability, reliability, and optimized performance in distributed storage systems and data centers. In recent years, with the emergence of new storage devices, heterogeneous object-based storage system, such as a storage system with the co-existence of hard disk drives and solid state drives, have become increasingly attractive as they combine merits of different storage devices to deliver better promise. However, existing data replication schemes do not place data based on heterogeneous device characteristics as well as considering distinct data access patterns. In this paper, we introduce a novel data replication scheme PRS to achieve efficient data replication for heterogeneous storage systems. Different from traditional schemes, the PRS groups objects according to data access patterns and distributes replicas to heterogeneous devices with their features. It uses a pseudo random algorithm to optimize replica layout by considering storage device performance and capacity. The experimental results confirm that PRS is a highly efficient replication scheme for heterogeneous storage systems. Wei Xie 0017, Dong Dai 0001, Yong Chen 0001 |
CCGrid | 3 |
| 2017 | IOGP: An Incremental Online Graph Partitioning Algorithm for Distributed Graph DatabasesabstractGraphs have become increasingly important in many applications and domains such as querying relationships in social networks or managing rich metadata generated in scientific computing. Many of these use cases require high-performance distributed graph databases for serving continuous updates from clients and, at the same time, answering complex queries regarding the current graph. These operations in graph databases, also referred to as online transaction processing (OLTP) operations, have specific design and implementation requirements for graph partitioning algorithms. In this research, we argue it is necessary to consider the connectivity and the vertex degree changes during graph partitioning. Based on this idea, we designed an Incremental Online Graph Partitioning (IOGP) algorithm that responds accordingly to the incremental changes of vertex degree. IOGP helps achieve better locality, generate balanced partitions, and increase the parallelism for accessing high-degree vertices of the graph. Over both real-world and synthetic graphs, IOGP demonstrates as much as 2x better query performance with a less than 10% overhead when compared against state-of-the-art graph partitioning algorithms. Dong Dai 0001, Wei Zhang 0097, Yong Chen 0001 |
HPDC | 1 |
| 2017 | POSTER: IOGP: An Incremental Online Graph Partitioning for Large-Scale Distributed Graph DatabasesabstractLarge-scale graphs are becoming critical in various domains such as social network, scientific application, knowledge discovery, and even system software, etc. Many of those use cases require large-scale high-performance graph databases, which are designed for serving continuous updates from the clients, and at the same time, answering complex queries towards the current graph in an on-line manner. Those operations in graph databases, also referred as OLTP (online transaction processing) operations, need specific design and implementation in graph partitioning algorithms. In this study, we designed an incremental online graph partitioning (IOGP), optimized for OLTP workloads. It is designed to achieve better locality, generate balanced partitions, and increase the parallelism for accessing hotspots of the graph. Our evaluation results on both real world and synthetic graphs in both simulation and real system confirm a better performance on graph queries (as much as 2X) with small overheads during graph insertion (less than 10%). Dong Dai 0001, Wei Zhang 0097, Yong Chen 0001 |
PPoPP | 1 |
| 2017 | SuperMIC: Analyzing Large Biological Datasets in Bioinformatics with Maximal Information CoefficientabstractThe maximal information coefficient (MIC) has been proposed to discover relationships and associations between pairs of variables. It poses significant challenges for bioinformatics scientists to accelerate the MIC calculation, especially in genome sequencing and biological annotations. In this paper, we explore a parallel approach which uses MapReduce framework to improve the computing efficiency and throughput of the MIC computation. The acceleration system includes biological data storage on HDFS, preprocessing algorithms, distributed memory cache mechanism, and the partition of MapReduce jobs. Based on the acceleration approach, we extend the traditional two-variable algorithm to multiple variables algorithm. The experimental results show that our parallel solution provides a linear speedup comparing with original algorithm without affecting the correctness and sensitivity. Chao Wang 0003, Dong Dai 0001, Xi Li 0003, Aili Wang 0003, Xuehai Zhou |
IEEE ACM Trans. Comput. Biol. Bioinform. | 2 |
| 2016 | GraphMeta: A Graph-Based Engine for Managing Large-Scale HPC Rich MetadataabstractHigh-performance computing (HPC) systems face increasingly critical metadata management challenges, especially in the approaching exascale era. These challenges arise not only from exploding metadata volumes but also from increasingly diverse metadata, which contains data provenance and user-defined attributes in addition to traditional POSIX metadata. This "rich" metadata is critical to support many advanced data management functionality such as data auditing and validation. In our prior work, we presented a graph-based model that could be a promising solution to uniformly manage such rich metadata because of its flexibility and generality. At the same time, however, graph-based rich metadata management introduces significant challenges. In this study, we first identify the challenges presented by the underlying infrastructure in supporting scalable, high-performance rich metadata management. To tackle these challenges, we then present GraphMeta, a graph-based engine designed for managing large-scale rich metadata. We also utilize a series of optimizations designed for rich metadata graphs. We evaluate GraphMeta with both synthetic and real HPC metadata workloads and compare it with other approaches. The results show that its advantages in terms of rich metadata management in HPC systems, including better performance and scalability compared with existing solutions. Dong Dai 0001, Yong Chen 0001, Philip H. Carns, John Jenkins, Wei Zhang 0097, Robert B. Ross |
CLUSTER | 1 |
| 2016 | An asynchronous traversal engine for graph-based rich metadata management
Dong Dai 0001, Philip H. Carns, Robert B. Ross, John Jenkins, Nicholas Muirhead, Yong Chen 0001 |
Parallel Comput. | 1 |
| 2015 | GraphTrek: Asynchronous Graph Traversal for Property Graph-Based Metadata ManagementabstractProperty graphs are a promising data model for rich metadata management in high-performance computing (HPC) systems because of their ability to represent not only metadata attributes but also the relationships between them. A property graph can be used to record the relationships between users, jobs, and data, for example, with unique annotations for each entity. This high-volume, power-law distributed use case is a natural fit for an out-of-core distributed property graph database. Such a system must support live updates (to ingest production information in real time), low-latency point queries (for frequent metadata operations such as permission checking), and large-scale traversals (for provenance data mining). Large-scale property graph traversals are particularly challenging for distributed graph databases, however. Most existing graph databases implement a "level-synchronous" breadth-first search algorithm that relies on global synchronization in each traversal step. This traversal model performs well in many problem domains, but a rich metadata management system is characterized by imbalanced graphs, long traversal lengths, and concurrent workloads, each of which has the potential to introduce or exacerbate stragglers. We define stragglers as abnormally slow steps (or servers) in a graph traversal that lead to low overall throughput for synchronous traversal algorithms. The straggler problem can be mitigated by the use of asynchronous traversal algorithms. Asynchronous traversal has been successfully demonstrated in graph processing frameworks, but such systems require the graph to be loaded into a separate batch-processing framework. In this work, we propose GraphTrek, a general asynchronous graph traversal engine working with graph databases for processing rich metadata management in their native format. We also outline a traversal-aware query language and key optimizations (traversal-affiliate caching and execution merging) necessary for efficient performance. Our experiments show that the asynchronous graph traversal engine is more efficient than its synchronous counterpart in the case of HPC rich metadata processing, where more servers are involved and larger traversals are needed. Dong Dai 0001, Philip H. Carns, Robert B. Ross, John Jenkins, Kyle Blauer, Yong Chen 0001 |
CLUSTER | 1 |
| 2014 | Provenance-based object storage prediction scheme for scientific big data applicationsabstractObject storage has been increasingly adopted in high-performance computing for scientific, big data applications. With object storage, applications usually use object IDs, queries, or collections to identify the data instead of using files. Since the object store changes the way data is accessed in applications, it introduces new challenges for I/O prediction, which used to work based on interfile or intrafile pattern detection. The key challenge is that the inputs of object-based applications are no longer expressed as static file names: they become much more dynamic and unstable, hidden inside application logic. Traditional prediction strategies do not work well in such conditions. In this paper, we introduce the use of provenance information, which was collected for data management in high-performance computing systems, in order to build an accurate coarse-grained (object-level) input prediction. The prediction results can be preloaded into a burst buffer to accelerate future reads. To our best knowledge, this study is the first to use provenance information in object stores to predict application inputs. Evaluation results confirm the effectiveness and accuracy of our provenance-based prediction and show that the proposed prediction system is feasible for real-work deployment. Dong Dai 0001, Yong Chen 0001, Dries Kimpe, Robert B. Ross |
IEEE BigData | 1 |
| 2014 | Provenance-Based Prediction Scheme for Object Storage System in HPCabstractObject-based storage model is recently widely adopted both in industry and academia to support growingly data intensive applications in high-performance computing. However, the I/O prediction strategies which have been proven effective in traditional parallel file systems, have not been thoroughly studied under this new object-based storage model. There are new challenges introduced from object storage that make traditional prediction systems not work properly. In this paper, we propose a new I/O access prediction system based on provenance analysis on both applications and objects. We argue that the provenance, which contains metadata that describes the history of data, reveals the detailed information about applications and data sets, which can be used to capture the system status and provide accurate I/O prediction efficiently. Our current evaluations based on real-world trace data (Darshan datasets) simulation also confirm that provenance-based prediction system is able to provide accurate predictions for object storage systems. Dong Dai 0001, Yong Chen 0001, Dries Kimpe, Robert B. Ross |
CCGRID | 1 |
| 2014 | Domino: an incremental computing framework in cloud with eventual synchronizationabstractIn recent years, more and more applications in cloud have needed to process large-scale on-line data sets that evolve over time as entries are added or modified. Several programming frameworks, such as Percolator and Oolong, are proposed for such incremental data processing and can achieve efficient updates with an event-driven abstraction. However, these frameworks are inherently asynchronous, leaving the heavy burden of managing synchronization to applications developers. Such a limitation significantly restricts their usability. In this paper, we introduce a trigger-based incremental computing framework, called Domino, with a flexible synchronization mechanism and runtime optimizations to coordinate parallel triggers efficiently. With this new framework, both synchronous and asynchronous applications can be seamlessly developed. Use cases and current evaluation results confirm that the new Domino programming model delivers sufficient performance and is easy to use in large-scale distributed computing. Dong Dai 0001, Yong Chen 0001, Dries Kimpe, Robert B. Ross, Xuehai Zhou |
HPDC | 1 |
| 2014 | Temperature-Aware Scheduling Based on Dynamic Time-Slice Scaling
Gangyong Jia, Youwei Yuan, Jian Wan 0001, Congfeng Jiang, Xi Li 0003, Dong Dai 0001 |
ICA3PP (1) | 6 |
| 2014 | An Adaptive Auto-configuration Tool for HadoopabstractWith the coming concept of 'big data', the ability to handle large datasets has become a critical consideration for the success of industrial organizations such as Google, Amazon, Yahoo! and Facebook. As an important Cloud Computing framework for bulk data processing, Hadoop is widely used in these organizations. However, the performance of MapReduce is seriously limited by its stiff configuration strategy. Even for a single simple job in Hadoop, a large number of tuning parameters have to be set by users. This may easily lead to performance loss due to some misconfigurations. In this paper, we present an adaptive automatic configuration tool (AACT) for Hadoop to achieve performance optimization. To achieve this goal, we propose a mathematical model which will accurately learn the relationship between system performance and configuration parameters, then configure Hadoop system based on this mathematical model. With the help of AACT, Hadoop is able to adapt the hardware and software configurations dynamically and drive the system to an optimal configuration in acceptable time. Experimental results show its efficiency and adaptability, and that it is ten times faster compared with default configuration. Hang Zhuang, Kun Lu 0002, Jinhong Zhou, Dong Dai 0001, Xuehai Zhou |
ICECCS | 6 |
| 2014 | Combine thread with memory scheduling for maximizing performance in multi-core systemsabstractThe growing gap between microprocessor speed and DRAM speed is a major problem that computer designers are facing. In order to narrow the gap, it is necessary to improve DRAM's speed and throughput. Moreover, on multi-core platforms, DRAM memory shared by all cores usually suffers from the memory contention and interference problem, which can cause serious performance degradation and unfairness among parallel running threads. To address these problems, this paper proposes techniques to take both advantages of partitioning cores, threads and memory banks into groups to reduce interference among different groups and grouping the memory accesses of the same row together to reduce cache miss rate. A memory optimization framework combined thread scheduling with memory scheduling (CTMS) is proposed in this paper, which simultaneously minimizes memory access schedule length, memory access time and reduce interference to maximize performance for multi-core systems. Experimental results show CTMS is 12.6% shorter in memory access time, while improving 11.8% throughput on average. Moreover, CTMS also saves 5.8% of the energy consumption. Gangyong Jia, Guangjie Han, Liang Shi 0001, Jian Wan 0001, Dong Dai 0001 |
ICPADS | 5 |
| 2014 | DLBS: Decentralized load balancing scheme for event-driven cloud frameworksabstractWith the development of cloud computing, more and more applications are moving to a distributed fashion to solve problems. These applications usually contain complex iterative or incremental procedures and have a more urgent requirement on low-latency. Thus many event-driven cloud frameworks are proposed. To optimize this kind of frameworks, an efficient strategy to minimize the execution time by redistributing work- loads is needed. Nowadays, load balance is a critical issue for the efficient operation of cloud platforms and many centralized schemes have already been proposed. However, few of them have been designed to support event-driven frameworks. Besides, as the cluster size and volume of tasks increases, centralized scheme will lead to a bottleneck of master node. In this paper, we demonstrate a decentralized load balancing scheme named DLBS for event-driven cloud frameworks and present two technologies to optimize it. In our design, schedulers are placed in every node for independently load-monitoring, autonomous decision-making and parallel task-scheduling. With the help of DLBS, master frees from the burden and tasks are executed with lower latency. We analyze the excellence of DLBS theoretically and proof it through simulation. At last, we implement and deploy it on a 64-machine cluster and demonstrate that it performs within 20% of an ideal scheme, which are consistent with simulation results. Xuehai Zhou, Kun Lu 0002, Jinhong Zhou, Hang Zhuang, Dong Dai 0001 |
ICPADS | 7 |
| 2014 | PUMA: Pseudo unified memory architecture for single-ISA heterogeneous multi-core systemsabstractSingle-ISA heterogeneous multi-core processors have advantages over cost-equivalent homogeneous ones, which integrate cores having the same instruction set architecture (ISA) but offer different performance and power characteristics. When these cores share the off-chip main memory, requests from different cores will interfere with each other, leading to low system performance and unfairness even starvation. Unfortunately, state-of-the-art memory scheduling and thread scheduling algorithms are ineffective at solving these problems. This paper proposes a fundamentally new memory architecture of pseudo unified memory (PUMA), which partitions the memory into regions according cores' different performance, each core mostly requests only one memory region seldom exceeding, reducing interfere among cores while retaining bank level parallelism for improving performance and fairness. We evaluate the design trade-offs involved in our PUMA and compare it against three state-of-the-art memory management methods. Our experimental results show that PUMA improves both system performance and fairness among cores while reducing memory power. Gangyong Jia, Liang Shi 0001, Jian Wan 0001, Youwei Yuan, Xi Li 0003, Dong Dai 0001 |
RTCSA | 6 |
| 2014 | Two-Choice Randomized Dynamic I/O Scheduler for Object Storage SystemsabstractObject storage is considered a promising solution for next-generation (exascale) high-performance computing platform because of its flexible and high-performance object interface. However, delivering high burst-write throughput is still a critical challenge. Although deploying more storage servers can potentially provide higher throughput, it can be ineffective because the burst-write throughput can be limited by a small number of stragglers (storage servers that are occasionally slower than others). In this paper, we propose a two-choice randomized dynamic I/O scheduler that schedules the concurrent burst-write operations in a balanced way to avoid stragglers and hence achieve high throughput. The contributions in this study are threefold. First, we propose a two-choice randomized dynamic I/O scheduler with collaborative probe and preassign strategies. Second, we design and implement a redirect table and metadata maintainer to address the metadata management challenge introduced by dynamic I/O scheduling. Third, we evaluate the proposed scheduler with both simulation tests and experimental tests in an HPC cluster. The evaluation results confirm the scalability and performance benefits of the proposed I/O scheduler. Dong Dai 0001, Yong Chen 0001, Dries Kimpe, Robert B. Ross |
SC | 1 |
| 2014 | Unbinds data and tasks to improving the Hadoop performanceabstractHadoop is a popular framework that provides easy programming interface of parallel programs to process large scale of data on clusters of commodity machines. Data intensive programs are the important part running on the cluster especially in large scale machine learning algorithm which executes of the same program iteratively. In-memory cache of input data is an efficient way to speed up these data intensive programs. However, we cannot be able to load all the data in memory because of the limitation of memory capacity. So, the key challenge is how we can accurately know when data should be cached in memory and when it ought to be released. The other problem is that memory capacity may even not enough to hold the input data of the running program. This leads to there is some data cannot be cached in memory. Prefetching is an effective method for such situation. We provide a unbinding technology which do not put the programs and data binded together before the real computation start. With unbinding technology, Hadoop can get a better performance when using caching and prefetching technology. We provide a Hadoop framework with unbinding technology named unbinding-Hadoop which decide the map tasks' input data in the map starting up phase, not at the job submission phase. Prefetching as well can be used in unbinding-Hadoop and can get better performance compared with the programs without unbinding. Evaluations on this system show that unbinding-Hadoop reduces the execution time of jobs by 40.2% and 29.2% with WordCount programs and K-means algorithm. Kun Lu 0002, Dong Dai 0001, Xuehai Zhou, Hang Zhuang |
SNPD | 2 |
| 2013 | HDFS+: Concurrent Writes Improvements for HDFSabstractHDFS is a popular distributed file system which provides high scalability and throughput. It lacks built-in support for multi-source data generating, which arise naturally in many applications including log mining, data analysis etc. There needs a data collection step before analysis in basic HDFS environment because of many data are in local disk, such as log. We proposed a solution which can compose many existent files to a single file and it is suitable for concurrent writes by many data producers. Programs only have to implements data processing against one single file without a data collection step when data analysis. We implemented HDFS+ by modifying existent HDFS, and evaluated with applications including log analysis. Our results show great throughput improvements in data concurrent writes. HDFS+ vastly simplifies the data collecting steps in data analysis procedure. Kun Lu 0002, Dong Dai 0001 |
CCGRID | 2 |
| 2013 | Coordinate Task and Memory Management for Improving Power Efficiency
Gangyong Jia, Xi Li 0003, Jian Wan 0001, Chao Wang 0003, Dong Dai 0001, Congfeng Jiang |
ICA3PP (1) | 5 |
| 2012 | Cloud Based Short Read Mapping ServiceabstractBioinformatics is an emerging field with seemingly limitless possibilities for advances in numerous scientific research and applications domains. In this paper, we summaries the explosive cutting-edge acceleration engines for the emerging short read mapping problems. What's more, we propose a novel Cloud based web service solution to the short read mapping problem in DNA sequencing, which greatly accelerates the task of aligning continuous incoming short length reads to uncertain known reference genomes. This approach is based on the pre-process of the reference genomes and iterative MapReduce jobs for aligning the continuous incoming reads. The MapReduce-based read-mapping algorithm is modeled after RMAP. Preliminary experimental results on incorporated MapReduce programming framework demonstrate that our proposed architecture and methods efficiently reduces the waiting time for large scale short reads applications. This architecture would be much important and efficient in future commercial personal gnome sequencing service. Dong Dai 0001, Xi Li 0003, Chao Wang 0003, Xuehai Zhou |
CLUSTER | 1 |
| 2012 | Phase Detection for Loop-Based Programs on Multicore ArchitecturesabstractPhase detection and behavior analysis have been major concerned to improve the performance as well as the system throughputs. However, for the distributed acceleration engines, the execution among different phases is much more difficult to be analyzed, especially for the loop based programs. With respect to the tasks in different iterations, how to efficiently detect the phases belonging to the same loop iteration or even across iterations is posing significant challenge. In this paper we propose a phase detection method for loop-based programs on multiprocessor system-on-chip (MPSoC). A cross compiling tool based on state-of-the-art ARM RVDS is employed to locate the hot spot function of the program. Based on the hot spots, we target the function optimization on a hadoop cluster for performance evaluation. The preliminary experimental results demonstrate that our proposed techniques can extract the hot block function with high accuracy and modest overheads. The method can be applied to guide the optimization and adaptive mapping scheme on MPSoC architectures. Chao Wang 0003, Xi Li 0003, Dong Dai 0001, Gangyong Jia, Xuehai Zhou |
CLUSTER | 3 |