EDBT 2026 Demo / reviewers in the wild / expert
Dongfang Zhao 0001
dblp:25/5030-1
· DBLP profile ↗
56ranked-venue papers
15as first author
21since 2021 · last 2026
0000-0002-0677-634XORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 23 · 8 first-author · 7 since 2021Artificial intelligence and machine learning · 16 · 4 first-author · 4 since 2021Databases, data management, data science and information retrieval · 16 · 3 first-author · 6 since 2021Applied, interdisciplinary, general and emerging computing · 15 · 5 first-author · 2 since 2021Graphics, computer vision, multimedia, augmented reality and games · 5 · 3 since 2021Software engineering, systems software and programming languages · 3 · 1 first-author · 1 since 2021Computer networks · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Order-Preserving Dimension Reduction for Multimodal Semantic EmbeddingabstractSearching for the k-nearest neighbors in multimodal data retrieval is computationally expensive, particularly due to the inherent difficulty in comparing similarity measures across different modalities. Recent advances in multimodal machine learning address this issue by mapping data into a shared embedding space; however, the high dimensionality of these embeddings (hundreds to thousands of dimensions) presents a challenge for time-sensitive vision applications. This work proposes Order-Preserving Dimension Reduction (OPDR), aiming to reduce the dimensionality of embeddings while preserving the ranking of KNN in the lower-dimensional space. One notable component of OPDR is a new measure function to quantify KNN quality as a global metric, based on which we derive a closed-form map between target dimensionality and key contextual parameters. We have integrated OPDR with multiple state-of-the-art dimension-reduction techniques, distance functions, and embedding models; experiments on a variety of multimodal datasets demonstrate that OPDR effectively retains recall high accuracy while significantly reducing computational costs. Chengyu Gong, Gefei Shen, Luanzheng Guo, Nathan R. Tallent, Dongfang Zhao 0001 |
AAAI | 5 |
| 2026 | SIVF: GPU-Resident IVF Index for Streaming Vector AnalyticsabstractGPU-accelerated Inverted File (IVF) index is one of the industry standards for large-scale vector search but relies on static VRAM layouts that hinder real-time mutability. Our benchmark and analysis reveal that existing designs of GPU IVF necessitate expensive CPU-GPU data transfers for index updates, causing system latency to spike from milliseconds to seconds in streaming scenarios. We present SIVF, a GPU-native index that enables high-velocity, in-place mutation via a series of new data structures and algorithms, such as conflict-free slab allocation and coalesced search on non-contiguous memory. SIVF has been implemented and integrated into the open-source vector search library, Faiss. Evaluation against baselines with diverse vector datasets demonstrates that SIVF reduces deletion latency by orders of magnitude compared to the state-of-the-arts. Furthermore, distributed experiments on a 12-GPU cluster demonstrate that SIVF exhibits near perfect linear scalability, achieving an aggregate ingestion throughput of 4.07 million vectors/s and a deletion throughput of 108.5 million vectors/s. Dongfang Zhao 0001 |
HPDC | 1 |
| 2026 | QPAD: Quantile-Preserving Approximate Dimension Reduction for Nearest Neighbors Preservation in High-Dimensional Vector Search
Jiuzhou Fu, Dongfang Zhao 0001 |
ICDE | 2 |
| 2026 | Reliable Non-Leveled Homomorphic Encryption for Web ServicesabstractWith the ubiquitous deployment of web services, ensuring data confidentiality has become a challenging imperative. Fully Homomorphic Encryption (FHE) presents a powerful solution for processing encrypted data; however, its widespread adoption is severely constrained by two fundamental bottlenecks: substantial computational overhead and the absence of a built-in automatic error correction mechanism. These limitations render the deployment of FHE in real-world, complex network environments impractical. Baigang Chen, Dongfang Zhao 0001 |
WWW | 2 |
| 2025 | IPBA: Imperceptible Perturbation Backdoor Attack in Federated Self-Supervised LearningabstractFederated Self-Supervised Learning (FSSL) combines the advantages of decentralized modeling and unlabeled representation learning, serving as a cutting-edge paradigm with strong potential for scalability and privacy preservation. Although FSSL has garnered increasing attention, research indicates that it remains vulnerable to backdoor attacks. Existing methods generally rely on visually obvious triggers, which makes it difficult to meet the requirements for stealth and practicality in real-world deployment. In this paper, we propose an imperceptible and effective backdoor attack method against FSSL, called IPBA. Our empirical study reveals that existing imperceptible triggers face a series of challenges in FSSL, particularly limited transferability, feature entanglement with augmented samples, and out-of-distribution properties. These issues collectively undermine the effectiveness and stealthiness of traditional backdoor attacks in FSSL. To overcome these challenges, IPBA decouples the feature distributions of backdoor and augmented samples, and introduces Sliced-Wasserstein distance to mitigate the out-of-distribution properties of backdoor samples, thereby optimizing the trigger generation process. Our experimental results on several FSSL scenarios and datasets show that IPBA significantly outperforms existing backdoor attack methods in performance and exhibits strong robustness under various defense mechanisms. Jiayao Wang 0004, Zhendong Zhao, Junwu Zhu, Dongfang Zhao 0001 |
ECAI | 7 |
| 2025 | GraphQL vs. REST: Investigating Performance and Scalability for Serverless Data PersistenceabstractServerless computing simplifies application deployment by removing the need for infrastructure management, with REST APIs being the common interface. Data persistence is essential for serverless applications because serverless functions are stateless and short-lived. To retain information across function executions, relational databases are commonly used to persist data, to ensure continuity, scalability, and reliable management. However, REST can lead to inefficiencies such as data over-fetching and under-fetching, which impact performance. GraphQL is a data query and manipulation language that supports client queries that specify what data is to be retrieved or modified. GraphQL resolvers fetch and transform data as specified by client queries.This paper compares the performance and scalability of GraphQL APIs as a database interface middleware in contrast to REST APIs implemented using serverless functions. We compare a GraphQL API implemented using the managed GraphQL AWS AppSync service, an unmanaged GraphQL API hosted using Apollo Server, and a traditional REST API implemented via the Amazon API Gateway and AWS Lambda. We compare these alternatives using clients implemented with serverless AWS Lambda functions, and also a local machine, an Amazon virtual machine (VM), and a Google VM. Our data APIs accessed a managed Amazon Aurora PostgreSQL cluster populated with the U.S. Centers for Medicare & Medicaid Services (CMS) Open Payments dataset. Our GraphQL implementations show content-dependent performance compared to REST, with Apollo Server demonstrating 25-67% faster average round-trip times vs. REST for most operations, but worse scalability than REST with very high concurrent workloads. Our findings provide practical guidance for developers selecting serverless architectures for GraphQL and REST APIs based on specific application requirements, network conditions, and expected request volumes. Runjie Jin, Robert Cordingly, Dongfang Zhao 0001, Wes Lloyd |
IC2E | 3 |
| 2025 | Distributed FaaSRunner: Enabling Reproducible Multi-node, Multi-threaded Function-as-a-Service Endpoint TestingabstractToday’s cloud native applications are often built using a service-oriented architecture supported by many microservices hosted using serverless Function-as-a-Service (FaaS) platforms. The vast majority of serverless function benchmarks and tests are performed using a single-client machine to generate workloads against highly scalable serverless backends. However, single-client test engines can lack sufficient computational resources and network bandwidth to adequately stress serverless backends. Testing tools such as Apache JMeter support orchestrating tests using multiple client nodes, but lack the ability to orchestrate sophisticated tests with distributed workload patterns common with real-world serverless workloads. In this paper, we introduce Distributed FaaSRunner, a distributed test tool which supports the ability to reproduce multi-node serverless workloads using traces derived by ingesting serverless function log files and randomly generated workload traces. We test Distributed FaaSRunner’s ability to precisely reproduce serverless function request dispatch and arrival time using various test cluster configurations. Using globally distributed clients, we predict latency to adjust workload trace event dispatch times to reproduce original request arrival latency. We demonstrate that Distributed FaaSRunner can reproduce both temporal and spatial characteristics of serverless workloads, enabling new capabilities to assess performance of FaaS platforms beyond traditional load testing. Tomoki Kondo, Austin Bomhold, Robert Cordingly, Dongfang Zhao 0001, Wes Lloyd |
IC2E | 4 |
| 2025 | ProHD: Projection-Based Hausdorff Distance ApproximationabstractThe Hausdorff distance (HD) is a robust measure of set dissimilarity, but computing it exactly on large, high-dimensional datasets is prohibitively expensive. We propose ProHD, a projection-guided approximation algorithm that dramatically accelerates HD computation while maintaining high accuracy. ProHD identifies a small subset of candidate “extreme” points by projecting the data onto a few informative directions (such as the centroid axis and top principal components) and computing the HD on this subset. This approach guarantees an underestimate of the true HD with a bounded additive error and typically achieves results within a few percent of the exact value. In extensive experiments on image, physics, and synthetic datasets (up to two million points in D = 256), ProHD runs 10-100× faster than exact algorithms while attaining 5-20× lower error than random sampling-based approximations. Our method enables practical HD calculations in scenarios like large vector databases and streaming data, where quick and reliable set distance estimation is needed. Jiuzhou Fu, Luanzheng Guo, Nathan R. Tallent, Dongfang Zhao 0001 |
ICDM | 4 |
| 2025 | ZTP: A Scalable and Lightweight Privacy-Preserving Blockchain via Scale-Free Quorums and Geometric FragmentationabstractEnsuring data privacy in blockchain systems remains challenging due to the heavy computational and communication costs of traditional cryptographic mechanisms. Existing solutions often suffer from limited scalability, high resource consumption, and inefficient tamper-proof key management. To address these challenges, we propose Zero Trust Privacy (ZTP), a lightweight framework for scalable on-chain privacy and secure distributed key management. ZTP introduces a hybrid quorum protocol using dynamic scale-free graph adjustments and a parallel data and key management mechanism based on the Geometric Fragmentation Technique (GFT), achieving efficient, tamper-resistant shard handling. To further enhance scalability, we incorporate a lightweight consensus protocol with parallel transaction processing, isolating transactions and key access from untrusted blockchain nodes. We implement and evaluate ZTP on a distributed blockchain prototype, demonstrating outstanding performance, achieving up to 49% fault tolerance, and delivering speedups of at least 55 × compared to state-of-the-art blockchain protocols. Our results highlight ZTP’s potential for resource-constrained and large-scale blockchain deployments. Abdullah Al-Mamun 0001, Dongfang Zhao 0001, Gagan Agrawal, Ahmed Aleroud, Mohamed I. Ibrahem |
ICPP | 2 |
| 2024 | Privacy-Preserving Artificial Intelligence on Edge Devices: A Homomorphic Encryption ApproachabstractRecent advancements in privacy-preserving artificial intelligence (AI) have paved the way for enhanced privacy in computational processes. A standing challenge, however, is the robust privacy preservation in AI algorithms, especially when integrated into edge devices and Internet-of-Thing (IoT) infrastructures. Most prevailing solutions have adopted traditional encryption methods which, though secure, often introduce significant overhead and potential dips in accuracy. In this study, we put forth an innovative approach, utilizing the CKKS encryption scheme, aiming to harmoniously balance computational efficiency with stringent data privacy. By harnessing the capabilities of Full Homomorphic Encryption (FHE) under the CKKS scheme, we ensure the preservation of privacy, successfully curbing the inherent noise traditionally linked with accuracy reductions in similar encryption-oriented solutions. Through comprehensive experiments, our approach showcased its potential as a strong contender for privacy preservation, demonstrating commendable performance across all tests, affirming that FHE is indeed viable for devices with constrained computational power and energy resources. Muhammad Jahanzeb Khan, Bo Fang 0002, Gaetano Cimino, Stefano Cirillo, Lei Yang 0001, Dongfang Zhao 0001 |
ICWS | 6 |
| 2024 | Towards Nonparametric Topological Layers in Neural Networks
Gefei Shen, Dongfang Zhao 0001 |
PAKDD (3) | 2 |
| 2024 | Domain adaptation of time series via contrastive learning with task-specific consistency
Qiushu Chen, Dongfang Zhao 0001, Linhua Jiang |
Appl. Intell. | 3 |
| 2023 | TopoCommit: A Topological Commit Protocol for Cross-Ledger Transactions in Scientific ComputingabstractWhile increasingly more applications are tempted to manage their data in decentralized systems, such as blockchains or distributed ledgers, the data exchange across multiple, potentially heterogeneous, decentralized systems remains an open problem: State-of-the-art protocols cannot meet one or more of the core requirements, such as atomicity, liveness, and scalability. Specifically, in the field of scientific computing, although a blockchain service was recently developed for scientific computing environments, the data exchanges and transactions among distinct ledgers are not supported. Observing that many modern scientific applications are collaborated on by multiple teams and the increasingly complicated (in-situ) workflows thereof, we argue that there is a pressing need to realize an efficient and scalable protocol for distinct ledgers to exchange data in scientific computing. This paper proposes a topological approach to enabling atomic, nonblocking, and scalable data exchanges among an arbitrary number of scientific ledgers in the context of collaborative scientific computing. Specifically, we construct a topological space formed by these ledgers—abstracting those nodes in a cross-ledger transaction as topological objects such as abstract simplex and simplicial complex. These topological objects, in turn, serve as the building blocks of a topological protocol, namely TopoCommit, under practical assumptions. We implement TopoCommit and integrate it into SciChain, a recently published distributed ledger for tracking scientific data provenance. The extensive evaluation of up to 1,008 nodes and 144 distinct ledgers on CloudLab shows that TopoCommit outperforms state-of-the-art protocols by up to 70×. Olamide Timothy Tawose, Lei Yang 0001, Dongfang Zhao 0001 |
CLUSTER | 3 |
| 2023 | SciLance: Mitigate Load Imbalance for Parallel Scientific Applications in Cloud EnvironmentsabstractElastic cloud computing provides new opportunities for accelerating the process of scientific discovery. However, unlike high-performance computing (HPC) systems that are built and optimized for workloads with intensive inter-node communication demands, the low-latency and high bandwidth communication capability is only enabled on a few very expensive high-end instance types in the cloud, which leads to poor cost-effectiveness. In addition, re-balancing the workload through extra data movement among compute nodes is a common way to mitigate the load imbalance issue in many scientific simulations, which further amplifies the communication pressure and makes it challenging to efficiently use cloud resources. To this end, we propose SciLance, which addresses the workload imbalance challenge by utilizing the heterogeneous and elastic resources offered by cloud platforms. Particularly, instead of moving data excessively among compute instances to balance the workload, SciLance dynamically adjusts the computer instances used for running parallel tasks based on the runtime imbalance identified through profiling. We prototype SciLance and perform extensive evaluation using adaptive mesh refinement (AMR) based scientific applications. The evaluation results demonstrate that SciLance can achieve up to 36.63% better performance with 16.91% lower cost for AMR-based simulation codes. Xinying Wang 0001, Lipeng Wan 0001, Scott Klasky, Dongfang Zhao 0001, Feng Yan 0001 |
CLUSTER | 4 |
| 2023 | Transparent Object Detection with Simulation Heatmap Guidance and Context Spatial Attention
Bobo Ju, Linhua Jiang, Dongfang Zhao 0001 |
MMM (2) | 5 |
| 2023 | Toward Efficient Homomorphic Encryption for Outsourced Databases through Parallel CachingabstractMany applications deployed to public clouds are concerned about the confidentiality of their outsourced data, such as financial services and electronic patient records. A plausible solution to this problem is homomorphic encryption (HE), which supports certain algebraic operations directly over the ciphertexts. The downside of HE schemes is their significant, if not prohibitive, performance overhead for data-intensive workloads that are very common for outsourced databases, or database-as-a-serve in cloud computing. The objective of this work is to mitigate the performance overhead incurred by the HE module in outsourced databases. To that end, this paper proposes a radix-based parallel caching optimization for accelerating the performance of homomorphic encryption (HE) of outsourced databases in cloud computing. The key insight of the proposed optimization is caching selected radix-ciphertexts in parallel without violating existing security guarantees of the primitive/base HE scheme. We design the radix HE algorithm and apply it to both batch- and incremental-HE schemes; we demonstrate the security of those radix-based HE schemes by showing that the problem of breaking them can be reduced to the problem of breaking their base HE schemes that are known IND-CPA (i.e. Indistinguishability under Chosen-Plaintext Attack). We implement the radix-based schemes as middleware of a 10-node Cassandra cluster on CloudLab; experiments on six workloads show that the proposed caching can boost state-of-the-art HE schemes, such as Paillier and Symmetria, by up to five orders of magnitude. Olamide Timothy Tawose, Jun Dai 0001, Lei Yang 0001, Dongfang Zhao 0001 |
Proc. ACM Manag. Data | 4 |
| 2022 | DEAN: A Lightweight and Resource-efficient Blockchain Protocol for Reliable Edge ComputingabstractEdge computing draws a lot of recent research interests because of the performance improvement by offloading many workloads from the remote data center to nearby edge nodes. Nonetheless, one open challenge of this emerging paradigm lies in the potential security issues on edge nodes. This paper proposes a cooperative protocol, namely DEAN, equipped with a unique resource-efficient quorum building mechanism to adopt blockchain seamlessly in an edge computing infrastructure to prevent data manipulation and allow fair data sharing with quick recovery under resource constraints of limited storage, computing, and network capacity. Specifically, DEAN leverages a parallel mechanism equipped with three independent core components, effectively achieving low resource consumption while allowing secured parallel block processing on edge nodes. We have implemented a system prototype based on DEAN and experimentally verified its effectiveness with a comparison with four popular blockchain implementations: Ethereum, Parity, IOTA, and Hyperledger Fabric. Experimental results show that the system prototype exhibits high resilience to arbitrary failures. Performance-wise, DEAN-based blockchain implementation out-performs the state-of-the-art blockchain systems with up to 88.6 x higher throughput and 26 x lower latency. Abdullah Al-Mamun 0001, Haoting Shen, Dongfang Zhao 0001 |
IPDPS | 3 |
| 2022 | Topological Modeling and Parallelization of Multidimensional Data on Microelectrode ArraysabstractMicroelectrode arrays (MEAs) are physical devices widely used in various science and engineering fields. One common computational challenge when applying a high-density MEA (i.e., a larger number of wires, more accurate locations of abnormal cells) is how to efficiently compute those resistance values provided the nonlinearity of the system of equations with the unknown resistance values per the Kirchhoff law. This paper proposes an algebraic-topological model for MEAs such that we can identify the intrinsic parallelism that cannot be identified by conventional approaches. We implement a system prototype called Parma based on the proposed topological methodology. Experimental results show that Parma outperforms the state-of-the-practice in time, scalability and memory usage: the computation time is two orders of magnitude faster on up to 1,024 cores with almost linear scalability and the memory is much better utilized with proportionally less warm-up time with respect to the number of concurrent threads. Olamide Timothy Tawose, Lei Yang 0001, Feng Yan 0001, Dongfang Zhao 0001 |
IPDPS | 5 |
| 2022 | Memory Scaling of Cloud-Based Big Data Systems: A Hybrid ApproachabstractWhen deploying applications with dynamic and intensive memory footprint to big data systems on public clouds, one important yet challenging question to answer is how to select a specific instance type whose memory capacity is large enough to prevent out-of-memory errors while the cost is minimized without violating performance requirements. The state-of-the-practice solution is trial and error, causing both performance overhead and additional monetary cost. This article investigates two memory scaling mechanisms in public clouds: physical memory (good performance and high cost) and virtual memory (degraded performance and no additional cost). In order to analyze the trade-off between performance and cost of the two scaling options, a performance-cost model is developed that is driven by a lightweight analytic prediction approach through a compact representation of the memory footprint. In addition, for those scenarios when the footprint is unavailable, a meta-model-based prediction method is proposed using just-in-time migration mechanisms. The proposed techniques have been extensively evaluated with various benchmarks and real-world applications on Amazon Web Services: the performance-cost model is highly accurate and the proposed just-in-time migration approach reduces the monetary cost by up to 66 percent. Xinying Wang 0001, Cong Xu 0002, Ke Wang 0012, Feng Yan 0001, Dongfang Zhao 0001 |
IEEE Trans. Big Data | 5 |
| 2021 | SciChain: Blockchain-enabled Lightweight and Efficient Data Provenance for Reproducible Scientific ComputingabstractThe state-of-the-art for auditing and reproducing scientific applications on high-performance computing (HPC) systems is through a data provenance subsystem. While recent advances in data provenance lie in reducing the performance overhead and improving the user's query flexibility, the fidelity of data provenance is often overlooked: there is no such way to ensure that the provenance data itself has not been fabricated or falsified. This paper advocates leveraging blockchains to deliver immutable and autonomous data provenance services such that scientific discoveries are trustworthy. The challenges for adopting blockchains to HPC include designing a new blockchain architecture compatible with the HPC platforms and, more importantly, a set of new consensus protocols for scientific applications atop blockchains. To this end, we have designed the proof-of-scalable-traceability (POST) protocol and implemented it in a blockchain prototype, namely SciChain, the very first practical blockchain system for provenance services on HPC. We evaluated SciChain by comparing it with multiple state-of-the-art systems; experimental results showed that SciChain guaranteed trustworthy data provenance while incurring orders of magnitude lower overhead than existing solutions. Abdullah Al-Mamun 0001, Feng Yan 0001, Dongfang Zhao 0001 |
ICDE | 3 |
| 2021 | BAASH: lightweight, efficient, and reliable blockchain-as-a-service for HPC systemsabstractDistributed resiliency becomes paramount to alleviate the growing costs of data movement and I/Os while preserving the data accuracy in HPC systems. This paper proposes to adopt blockchain-like decentralized protocols to achieve such distributed resiliency. The key challenge for such an adoption lies in the mismatch between blockchain's targeting systems (e.g., shared-nothing, loosely-coupled, TCP/IP stack) and HPC's unique design on storage subsystems, resource allocation, and programming models. We present BAASH, Blockchain-As-A-Service for HPC, deployable in a plug-n-play fashion. BAASH bridges the HPC-blockchain gap with two key components: (i) Lightweight consensus protocols for the HPC's shared-storage architecture, (ii) A new fault-tolerant mechanism compensating for the MPI to guarantee the distributed resiliency. We have implemented a prototype system and evaluated it with more than two million transactions on a 500-core HPC cluster. Results show that the prototype of the proposed techniques significantly outperforms vanilla blockchain systems and exhibits strong reliability with MPI. Abdullah Al-Mamun 0001, Feng Yan 0001, Dongfang Zhao 0001 |
SC | 3 |
| 2020 | HDK: Toward High-Performance Deep-Learning-Based Kirchhoff Analysis
Xinying Wang 0001, Olamide Timothy Tawose, Feng Yan 0001, Dongfang Zhao 0001 |
AAAI | 4 |
| 2020 | Reflector: a fine-grained I/O tracker for HPC systemsabstractWe present Reflector, to support both high-level and low-level I/O monitoring through user-defined interfaces such as HDF5 and NetCDF in addition to POSIX- and MPI-IO. We evaluate Reflector on both an on-premises 500-core HPC cluster and a leadership-class supercomputer at the Lawrence Berkeley National Laboratory. Preliminary results are promising as the system prototype incurs negligible performance overhead and clearly illustrates the I/O patterns and bottlenecks of multiple applications. Abdullah Al-Mamun 0001, Jialin Liu 0002, Tonglin Li, Quincey Koziol, Zhongyi Zhai, Junyan Qian, Haoting Shen, Dongfang Zhao 0001 |
PPoPP | 8 |
| 2020 | Locality-Aware Scheduling for Containers in Cloud ComputingabstractThe state-of-the-art scheduler of containerized cloud services considers load balance as the only criterion; many other important properties, including application performance, are overlooked. In the era of Big Data, however, applications evolve to be increasingly more data-intensive thus perform poorly when deployed on containerized cloud services. To that end, this paper aims to improve today's cloud service by taking application performance into account for the next-generation container schedulers. More specifically, in this work we build and analyze a new model that respects both load balance and application performance. Unlike prior studies, our model abstracts the dilemma between load balance and application performance into a unified optimization problem and then employs a statistical method to efficiently solve it. The most challenging part is that some sub-problems are extremely complex (for example, NP-hard), and heuristic algorithms have to be devised. Last but not least, we implement a system prototype of the proposed scheduling strategy for containerized cloud services. Experimental results show that our system can significantly boost application performance while preserving high load balance. Dongfang Zhao 0001, Mohamed Mohamed 0001, Heiko Ludwig |
IEEE Trans. Cloud Comput. | 1 |
| 2019 | Swift machine learning model serving scheduling: a region based reinforcement learning approachabstractThe success of machine learning has prospered Machine-Learning-as-a-Service (MLaaS) - deploying trained machine learning (ML) models in cloud to provide low latency inference services at scale. To meet latency Service-Level-Objective (SLO), judicious parallelization at both request and operation levels is utterly important. However, existing ML systems (e.g., Tensorflow) and cloud ML serving platforms (e.g., SageMaker) are SLO-agnostic and rely on users to manually configure the parallelism. To provide low latency ML serving, this paper proposes a swift machine learning serving scheduling framework with a novel Region-based Reinforcement Learning (RRL) approach. RRL can efficiently identify the optimal parallelism configuration under different workloads by estimating performance of similar configurations with that of the known ones. We both theoretically and experimentally show that the RRL approach can outperform state-of-the-art approaches by finding near optimal solutions over 8 times faster while reducing inference latency up to 79.0% and reducing SLO violation up to 49.9%. Heyang Qin, Syed Zawad, Yanqi Zhou, Lei Yang 0001, Dongfang Zhao 0001, Feng Yan 0001 |
SC | 5 |
| 2019 | Distributed Real-Time Data Aggregation Scheduling in Duty-Cycled Multi-hop Sensor Networks
Xiaohua Xu 0002, Yi Zhao 0004, Dongfang Zhao 0001, Lei Yang 0001, Spiridon Bakiras |
WASA | 3 |
| 2018 | Toward Cost-Effective Memory Scaling in Clouds: Symbiosis of Virtual and Physical MemoryabstractWhen deploying memory-intensive applications to public clouds, one important yet challenging problem is selecting a specific instance type whose memory capacity is large enough to prevent out-of-memory errors while the cost is minimized without violating performance requirements. The state-of-the-practice solution is trial and error, causing both performance overhead and additional monetary cost. This paper investigates two memory scaling mechanisms in public cloud: physical memory (good performance and high cost) and virtual memory (degraded performance and no additional cost). In order to analyze the trade-off between performance and cost of the two scaling options, a performance-cost model is developed that is driven by a lightweight analytic prediction approach through a compact representation of the memory footprint. In addition, for those scenarios when the footprint is unavailable, a meta-model based prediction method is proposed using just-in-time migration mechanisms. The proposed techniques have been extensively evaluated with various benchmarks and real-world applications on Amazon Web Services: the performance-cost model is highly accurate with errors ranging from 1% to 4% and the proposed just-in-time migration approach reduces the monetary cost by up to 66%. Xinying Wang 0001, Cong Xu 0002, Ke Wang 0012, Feng Yan 0001, Dongfang Zhao 0001 |
IEEE CLOUD | 5 |
| 2018 | In-memory Blockchain: Toward Efficient and Trustworthy Data Provenance for HPC SystemsabstractThe state-of-the-art approaches for tracking data provenance on high-performance computing (HPC) systems are either supported by file systems or relational databases. These techniques shared the same critique on the provenance data’s fidelity and the associated I/O overhead. This paper envisions to track the HPC data provenance using a distributed in-memory ledger—the core technique leveraged by blockchains and proven to be highly trustworthy by many large-scale applications. We pinpoint two system challenges—storage architecture and consensus protocol—for adopting blockchains to HPC and make the following contributions: (i) We design a new in-memory blockchain architecture for HPC systems, exploiting the high-performance network infrastructure InfiniBand and greatly reducing the I/O overhead; and (ii) We develop a new consensus protocol, namely proof-of-reproducibility (PoR), crafted for the new architecture, which takes into account both proof-of-work (PoW) and proof-of-stake (PoS) mechanisms. The correctness of PoR is both theoretically proven and experimentally verified. A prototype system is implemented and evaluated with more than one million transactions, showing 32× speedup compared to the filesystem-based provenance service and four orders of magnitude speedup compared to the database-based provenance service. Abdullah Al-Mamun 0001, Tonglin Li, Mohammad Sadoghi, Dongfang Zhao 0001 |
IEEE BigData | 4 |
| 2018 | Toward Scalable Analysis of Multidimensional Scientific Data: A Case Study of Electrode ArraysabstractMany modern scientific applications involve large volumes of multidimensional data and extensive computation. Although distributed systems and tools are becoming increasingly scalable, they are still far away to catch up the exponential growth rate exhibited by many of those scientific big-data applications. This paper presents our early effort on overcoming the exponential complexity of one widely deployed workload over multidimensional scientific data—the n×n numerical analysis on two-dimensional arrays. More specifically, we propose a new approach to reduce the exponentially-grown data into a semantically-equivalent polynomial form in the context of two-dimensional electrode arrays, which are widely used in biomedical engineering, electrical engineering, and mechanical engineering. We have implemented a system prototype in Python, preliminary results show that the proposed approach outperforms the state-of-the-practice in various metrics: (i) the consumed space is six orders of magnitude smaller; (ii) the execution time is three orders of magnitude faster; and (iii) the scalability is improved by two orders of magnitude—from 6×6 to 100 × 100—on mainstream servers in reasonable time. Ye Niu, Abdullah Al-Mamun 0001, Tonglin Li, Yi Zhao 0004, Dongfang Zhao 0001 |
IEEE BigData | 6 |
| 2018 | Davram: Distributed Virtual Memory in User SpaceabstractOne of the most challenging problems in modern distributed big data systems lies in their memory management: these systems preallocate a fixed amount of memory before applications start. In the best case where more memory can be acquired, users have to reconfigure the deployment and re-compute many intermediate results. If no more memory is available, users are then forced to manually partition the job into smaller tasks, incurring both development and performance overhead. This paper presents a user-level utility for managing the memory in distributed systems-the Distributed and Autonomous Virtual RAM (Davram). Davram enables to efficiently swap data between memory and disk in a distributed system without users' intervention or applications' awareness. Linhua Jiang, Ke Wang 0012, Dongfang Zhao 0001 |
CCGrid | 3 |
| 2018 | Toward Performant and Energy-efficient Queries in Three-tier Wireless Sensor NetworksabstractIn a wireless sensor network (WSN) where nodes are mostly battery-powered, queries' energy consumption and response time are two of the most important metrics as they represent the network's sustainability and performance, respectively. Conventional techniques used to focus only one of the two metrics and did not attempt to optimize both in a coordinated manner. This work aims to achieve both high sustainability and high performance of WSN queries at the same time. To that end, a new mechanism is proposed to construct the topology of a three-tier WSN. The proposed mechanism eliminates routing tables and employs a novel and efficient addressing scheme inspired by the Chinese Remainder Theorem (CRT). The CRT-based topology allows for query parallelism, an unprecedented feature in WSNs. On top of the new topology encoded by CRT, a new protocol is designed to parallelly preprocess collected data on sensor nodes by effectively aggregating and deduplicating data in a neighborhood cluster. Moreover, a new algorithm is devised to allow the queries and results to be transmitted through low-power and fault-tolerant paths using recursive elections over a subset of the entire power range. With all these new techniques combined, the proposed system outperforms the state-of-the-art from various perspectives: (i) the query response is improved by up to 53%; (ii) the energy consumption is reduced by up to 70%; and (iii) the reliability is increased by up to 39%. Jiayao Wang 0004, Abdullah Al-Mamun 0001, Tonglin Li, Linhua Jiang, Dongfang Zhao 0001 |
ICPP | 5 |
| 2017 | Comparative Evaluation of Big-Data Systems on Scientific Image Analytics WorkloadsabstractScientific discoveries are increasingly driven by analyzing large volumes of image data. Many new libraries and specialized database management systems (DBMSs) have emerged to support such tasks. It is unclear how well these systems support real-world image analysis use cases, and how performant the image analytics tasks implemented on top of such systems are. In this paper, we present the first comprehensive evaluation of large-scale image analysis systems using two real-world scientific image data processing use cases. We evaluate five representative systems (SciDB, Myria, Spark, Dask, and TensorFlow) and find that each of them has shortcomings that complicate implementation or hurt performance. Such shortcomings lead to new research opportunities in making large-scale image analysis both efficient and easy to use. Parmita Mehta, Sven Dorkenwald, Dongfang Zhao 0001, Tomer Kaftan, Alvin Cheung, Magdalena Balazinska, Ariel Rokem, Andrew J. Connolly, Jacob VanderPlas, Yusra AlSayyad |
Proc. VLDB Endow. | 3 |
| 2017 | Toward Efficient and Flexible Metadata Indexing of Big Data SystemsabstractIn Big Data era, applications are generating orders of magnitude more data in both volume and quantity. While many systems emerge to address such data explosion, the fact that these data's descriptors, i.e., metadata, are also “big” is often overlooked. The conventional approach to address the big metadata issue is to disperse metadata into multiple machines. However, it is extremely difficult to preserve both load-balance and data-locality in this approach. To this end, in this work we propose hierarchical indirection layers for indexing the underlying distributed metadata. By doing this, data locality is achieved efficiently by the indirection while load-balance is preserved. Three key challenges exist in this approach, however: first, how to achieve high resilience; second, how to ensure flexible granularity; third, how to restrain performance overhead. To address above challenges, we design Dindex, a distributed indexing service for metadata. Dindex incorporates a hierarchy of coarse-grained aggregation and horizontal key-coalition. Theoretical analysis shows that the overhead of building Dindex is compensated by only two or three queries. Dindex has been implemented by a lightweight distributed key-value store and integrated to a fully-fledged distributed filesystem. Experiments demonstrated that Dindex accelerated metadata queries by up to 60 percent with a negligible overhead. Dongfang Zhao 0001, Kan Qiao, Zhou Zhou 0006, Tonglin Li, Zhihan Lyu, Xiaohua Xu 0002 |
IEEE Trans. Big Data | 1 |
| 2016 | Toward Real-Time and Fine-Grained Monitoring of Software-Defined Networking in the CloudabstractPopular cloud infrastructure such as OpenStack, although incorporating a sophisticated module (e.g., Neutron) to handle its multi-layer network abstractions (from physical nodes, to virtual machines, to containers, all of which are multiplied by user-defined subnets), lacks the ability to allow users to both adjust network properties at runtime and track packet-wise activities. To this end, this position paper presents a system designed to enable real-time and fine-grained monitoring of the software-defined networking in cloud computing. While the full integration into OpenStack is still ongoing, we believe that the design, implementation, and open questions of our system development are insightful to the community. Dongfang Zhao 0001 |
CLOUD | 1 |
| 2016 | Albatross: An efficient cloud-enabled task scheduling and execution framework using distributed message queuesabstractData Analytics has become very popular on large datasets in different organizations. It is inevitable to use distributed resources such as Clouds for Data Analytics and other types of data processing at larger scales. To effectively utilize all system resources, an efficient scheduler is needed, but the traditional resource managers and job schedulers are centralized and designed for larger batch jobs which are fewer in number. Frameworks such as Hadoop and Spark, which are mainly designed for Big Data analytics, have been able to allow for more diversity in job types to some extent. However, even these systems have centralized architectures and will not be able to perform well on large scales and under heavy task loads. Modern applications generate tasks at very high rates that can cause significant slowdowns on these frameworks. Additionally, over-decomposition has shown to be very useful in increasing the system utilization. In order to achieve high efficiency, scalability, and better system utilization, it is critical for a modern scheduler to be able to handle over-decomposition and run highly granular tasks. Further, to achieve high performance, Albatross is written in C/C++, which imposes a minimal overhead to the workload process as compared to languages like Java or Python. We propose Albatross, a task level scheduling and execution framework that uses a Distributed Message Queue (DMQ) for task distribution among its workers. Unlike most scheduling systems, Albatross uses a pulling approach as opposed to the common push approach. The former would let Albatross achieve a good load balancing and scalability. Furthermore, the framework has built in support for task execution dependency on workflows. Therefore, Albatross is able to run various types of workloads, including Data Analytics and HPC applications. Finally, Albatross provides data locality support. This allows the framework to achieve higher performance through minimizing the amount of unnecessary data movement on the network. Our evaluations show that Albatross outperforms Spark and Hadoop at larger scales and in the case of running higher granularity workloads. Iman Sadooghi, Geet Kumar, Ke Wang 0012, Dongfang Zhao 0001, Tonglin Li, Ioan Raicu |
eScience | 4 |
| 2016 | A convergence of key-value storage systems from clouds to supercomputersabstractSummary This paper presents a convergence of distributed key‐value storage systems in clouds and supercomputers. It specifically presents ZHT, a zero‐hop distributed key‐value store system, which has been tuned for the requirements of high‐end computing systems. ZHT aims to be a building block for future distributed systems, such as parallel and distributed file systems, distributed job management systems, and parallel programming systems. ZHT has some important properties, such as being lightweight, dynamically allowing nodes join and leave, fault tolerant through replication, persistent, scalable, and supporting unconventional operations such as append, compare and swap, callback in addition to the traditional insert/lookup/remove. We have evaluated ZHT's performance under a variety of systems, ranging from a Linux cluster with 64 nodes, an Amazon EC2 virtual cluster up to 96 nodes, to an IBM Blue Gene/P supercomputer with 8K nodes. We compared ZHT against other key‐value stores and found it offers superior performance for the features and portability it supports. This paper also presents several real systems that have adopted ZHT, namely, FusionFS (a distributed file system), IStore (a storage system with erasure coding), MATRIX (distributed scheduling), Slurm++ (distributed HPC job launch), Fabriq (distributed message queue management); all of these real systems have been simplified because of key‐value storage systems and have been shown to outperform other leading systems by orders of magnitude in some cases. It is important to highlight that some of these systems are rooted in HPC systems from supercomputers, while others are rooted in clouds and ad hoc distributed systems; through our work, we have shown how versatile key‐value storage systems can be in such a variety of environments. Copyright © 2015 John Wiley & Sons, Ltd. Tonglin Li, Xiaobing Zhou, Ke Wang 0012, Dongfang Zhao 0001, Iman Sadooghi, Zhao Zhang 0007, Ioan Raicu |
Concurr. Comput. Pract. Exp. | 4 |
| 2016 | Exploiting multi-cores for efficient interchange of large messages in distributed systemsabstractSummary Conventional data serialization tools assume that objects to be coded are usually small in size so a single CPU core can encode it in a timely manner. In the era of Big Data, however, object gets increasingly complex and larger, which makes data serialization become a new performance bottleneck. This paper describes an approach to parallelize data serialization by leveraging multiple cores. Parallelizing data serialization introduces new questions such as how to split the (sub)objects, how to allocate the available cores, and how to minimize its overhead in practice. In this paper we design a framework for parallelly serializing large objects and analyze the design tradeoffs under different scenarios. To validate the proposed approach, we implemented parallel protocol buffers—the parallel version of Google's Protocol Buffers, a widely‐used data serialization utility. Experimental results confirm the effectiveness of Parallel Protocol Buffers: multiple cores employed in data serialization achieve highly scalable performance and incur negligible overhead. Copyright © 2015 John Wiley & Sons, Ltd. Dongfang Zhao 0001, Kan Qiao, Zhou Zhou 0006, Tonglin Li, Xiaobing Zhou, Ioan Raicu |
Concurr. Comput. Pract. Exp. | 1 |
| 2016 | Toward high-performance key-value stores through GPU encoding and locality-aware encoding
Dongfang Zhao 0001, Ke Wang 0012, Kan Qiao, Tonglin Li, Iman Sadooghi, Ioan Raicu |
J. Parallel Distributed Comput. | 1 |
| 2016 | I/O-aware bandwidth allocation for petascale computing systems
Zhou Zhou 0006, Xu Yang 0009, Dongfang Zhao 0001, Paul M. Rich, Wei Tang 0001, Zhiling Lan |
Parallel Comput. | 3 |
| 2016 | Towards Exploring Data-Intensive Scientific Applications at Extreme Scales through Systems and SimulationsabstractThe state-of-the-art storage architecture of high-performance computing systems was designed decades ago, and with today's scale and level of concurrency, it is showing significant limitations. Our recent work proposed a new architecture to address the I/O bottleneck of the conventional wisdom, and the system prototype (FusionFS) demonstrated its effectiveness on up to 16 K nodes-the scale on par with today's largest supercomputers. The main objective of this paper is to investigate FusionFS's scalability towards exascale. Exascale computers are predicted to emerge by 2018, comprising millions of cores and billions of threads. We built an event-driven simulator (FusionSim) according to the FusionFS architecture, and validated it with FusionFS's traces. FusionSim introduced less than 4 percent error between its simulation results and FusionFS traces. With FusionSim we simulated workloads on up to two million nodes and find out almost linear scalability of I/O performance; results justified FusionFS's viability for exascale systems. In addition to the simulation work, this paper extends the FusionFS system prototype in the following perspectives: (1) the fault tolerance of file metadata is supported, (2) the limitations of the current system design is discussed, and (3) a more thorough performance evaluation is conducted, such as N-to-1 metadata write, system efficiency, and more platforms such as Amazon Cloud. Dongfang Zhao 0001, Ning Liu 0008, Dries Kimpe, Robert B. Ross, Xian-He Sun, Ioan Raicu |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2016 | Dynamic Virtual Chunks: On Supporting Efficient Accesses to Compressed Scientific DataabstractData compression could ameliorate the I/O pressure of data-intensive scientific applications. Unfortunately, the conventional wisdom of naively applying data compression to the file or block brings the dilemma between efficient random accesses and high compression ratios. File-level compression barely supports efficient random accesses to the compressed data: any retrieval request need trigger the decompression from the beginning of the compressed file. Block-level compression provides flexible random accesses to the compressed blocks, but introduces extra overhead when applying the compressor to each and every block that results in a degraded overall compression ratio. This paper extends our prior work that introduces virtual chunks offering efficient random accesses to the compressed scientific data without sacrificing the compression ratio. Virtual chunks are logical blocks pointed at by appended references without breaking the physical continuity of the file content. These references allow the decompression to start from an arbitrary position (efficient random accesses), while no per-block overhead is introduced because the file's physical entirety is retained (high compression ratio). One limitation of virtual chunk is it only supports static references. This paper presents the algorithms, analysis, and evaluations of dynamic virtual chunks to deal with the cases where the references are updated dynamically. Dongfang Zhao 0001, Kan Qiao, Jian Yin 0002, Ioan Raicu |
IEEE Trans. Serv. Comput. | 1 |
| 2015 | A flexible QoS fortified distributed key-value storage system for the cloudabstractIn the era of big data and cloud, distributed key-value stores are increasingly used as building blocks of large-scale applications. Comparing to traditional relational databases, key-value stores are particularly compelling due to their low latency and excellent scalability. Many big companies, such as Facebook and Amazon, run multiple different applications and services on top of a single key-value store deployment to reduce the deployment and maintenance complexity as well as economic cost. However, every application has its performance requirement but most current key-value store systems are designed to serve every application request equally. This design works well when a single application accesses the key-value store, but it is not as good for the emerging concurrent multi-application scenario. In this paper, we present ZHT/Q, a flexible QoS (Quality of Service) fortified distributed key-value storage system for clouds and data centers. It improves the overall throughput by an order of magnitude and still satisfies different applications' latency requirements with QoS using dynamic and adaptive request batching mechanisms. The experiment results show that our new system delivers up to 28 times higher throughput than the base solution while more than 99% of requests' latency requirements are satisfied. Tonglin Li, Ke Wang 0012, Dongfang Zhao 0001, Kan Qiao, Iman Sadooghi, Xiaobing Zhou, Ioan Raicu |
IEEE BigData | 3 |
| 2015 | Data confidentiality challenges in big data applicationsabstractIn this paper, we address the problem of data confidentiality in big data analytics. In many fields, much useful patterns can be extracted by applying machine learning techniques to big data. However, data confidentiality must be protected. In many scenarios, data confidentiality could well be a prerequisite for data to be shared. We present a scheme to provide provable secure data confidentiality and discuss various techniques to optimize performance of such a system. Jian Yin 0002, Dongfang Zhao 0001 |
IEEE BigData | 2 |
| 2015 | Toward locality-aware scheduling for containerized cloud servicesabstractThe state-of-the-art scheduler of containerized cloud services considers load-balance as the only criterion and neglects many others such as application performance. In the era of Big Data, however, applications have evolved to be highly data-intensive thus perform poorly in existing systems. This particularly holds for Platform-as-a-Service environments that encourage an application model of stateless application instances in containers reading and writing data to services storing states, e.g., key-value stores. To this end, this work strives to improve today's cloud services by incorporating sensitivity to both load-balance and application performance. We built and analyzed theoretical models that respect both dimensions, and unlike prior studies, our model abstracts the dilemma between load-balance and application performance into an optimization problem and employs a statistical method to meet the discrepant requirements. Using heuristic algorithms and approaches we try to solve the abstracted problems. We implemented the proposed approach in Diego (an open-source cloud service scheduler) and demonstrate that it can significantly boost the performance of containerized applications while preserving a relatively high load-balance. Dongfang Zhao 0001, NagaPramod Mandagere, Gabriel Alatorre, Mohamed Mohamed 0001, Heiko Ludwig |
IEEE BigData | 1 |
| 2015 | MHT: A light-weight scalable zero-hop MPI enabled distributed key-value storeabstractIn this paper, we propose and implement a key-value store that supports MPI while allowing application access at any time without having to declaring in the same MPI communication world. This feature may significantly simplify the application design and allow programmers leverage the power of key-value store in an intuitive way. In our preliminary experiment results captured from a supercomputer at Los Alamos National Laboratory, our prototype shows linear scalability at up to 256 nodes. Xiaobing Zhou, Tonglin Li, Ke Wang 0012, Dongfang Zhao 0001, Iman Sadooghi, Ioan Raicu |
IEEE BigData | 4 |
| 2015 | GRAPH/Z: A Key-Value Store Based Scalable Graph Processing SystemabstractThe emerging applications in big data and social networks issue rapidly increasing demands on graph processing. Graph query operations that involve a large number of vertices and edges can be tremendously slow on traditional databases. The state-of-the-art graph processing systems and databases usually adopt master/slave architecture that potentially impairs their The contributions of this paper are as follows: scalability. This work describes the design and implementation of a new graph processing system based on Bulk Synchronous Parallel model. Our system is built on top of ZHT, a scalable distributed key-value store, which benefits the graph processing in terms of scalability, performance and persistency. The experiment results imply excellent scalability. Tonglin Li, Chaoqi Ma, Xiaobing Zhou, Ke Wang 0012, Dongfang Zhao 0001, Iman Sadooghi, Ioan Raicu |
CLUSTER | 6 |
| 2015 | I/O-Aware Batch Scheduling for Petascale Computing SystemsabstractIn the Big Data era, the gap between the storage performance and an application's I/O requirement is increasing. I/O congestion caused by concurrent storage accesses from multiple applications is inevitable and severely harms the performance. Conventional approaches either focus on optimizing an application's access pattern individually or handle I/O requests on a low-level storage layer without any knowledge from the upper-level applications. In this paper, we present a novel I/O-aware batch scheduling framework to coordinate ongoing I/O requests on petascale computing systems. The motivation behind this innovation is that the batch scheduler has a holistic view of both the system state and jobs' activities and can control the jobs' status on the fly during their execution. We treat a job's I/O requests as periodical subjobs within its lifecycle and transform the I/O congestion issue into a classical scheduling problem. We design two scheduling polices with different scheduling objectives either on user-oriented metrics or system performance. We conduct extensive trace-based simulations using real job traces and I/O traces from a production IBM Blue Gene/Q system. Experimental results demonstrate that our design can improve job performance by more than 30%, as well as increasing system performance. Zhou Zhou 0006, Xu Yang 0009, Dongfang Zhao 0001, Paul M. Rich, Wei Tang 0001, Zhiling Lan |
CLUSTER | 3 |
| 2014 | Optimizing load balancing and data-locality with data-aware schedulingabstractLoad balancing techniques (e.g. work stealing) are important to obtain the best performance for distributed task scheduling systems that have multiple schedulers making scheduling decisions. In work stealing, tasks are randomly migrated from heavy-loaded schedulers to idle ones. However, for data-intensive applications where tasks are dependent and task execution involves processing a large amount of data, migrating tasks blindly yields poor data-locality and incurs significant data-transferring overhead. This work improves work stealing by using both dedicated and shared queues. Tasks are organized in queues based on task data size and location. We implement our technique in MATRIX, a distributed task scheduler for many-task computing. We leverage distributed key-value store to organize and scale the task metadata, task dependency, and data-locality. We evaluate the improved work stealing technique with both applications and micro-benchmarks structured as direct acyclic graphs. Results show that the proposed data-aware work stealing technique performs well. Ke Wang 0012, Xiaobing Zhou, Tonglin Li, Dongfang Zhao 0001, Michael Lang 0003, Ioan Raicu |
IEEE BigData | 4 |
| 2014 | Virtual chunks: On supporting random accesses to scientific data in compressible storage systemsabstractData compression could ameliorate the I/O pressure of scientific applications on high-performance computing systems. Unfortunately, the conventional wisdom of naively applying data compression to the file or block brings the dilemma between efficient random accesses and high compression ratios. Filelevel compression can barely support efficient random accesses to the compressed data: any retrieval request need trigger the decompression from the beginning of the compressed file. Block-level compression provides flexible random accesses to the compressed data, but introduces extra overhead when applying the compressor to each every block that results in a degraded overall compression ratio. This paper introduces a concept called virtual chunks aiming to support efficient random accesses to the compressed scientific data without sacrificing its compression ratio. In essence, virtual chunks are logical blocks identified by appended references without breaking the physical continuity of the file content. These additional references allow the decompression to start from an arbitrary position (efficient random access), and retain the file's physical entirety to achieve high compression ratio on par with file-level compression. One potential concern of virtual chunks lies on its space overhead (from the additional references) that degrades the compression ratio, but our analytic study and experimental results demonstrate that such overhead is negligible. We have implemented virtual chunks in two forms: a middleware to the GPFS parallel file system, and a module in the FusionFS distributed file system. Large-scale evaluations on up to 1,024 cores showed that virtual chunks could help improve the I/O throughput by 2X speedup. Dongfang Zhao 0001, Jian Yin 0002, Kan Qiao, Ioan Raicu |
IEEE BigData | 1 |
| 2014 | FusionFS: Toward supporting data-intensive scientific applications on extreme-scale high-performance computing systemsabstractState-of-the-art, yet decades-old, architecture of high-performance computing systems has its compute and storage resources separated. It thus is limited for modern data-intensive scientific applications because every I/O needs to be transferred via the network between the compute and storage resources. In this paper we propose an architecture that hss a distributed storage layer local to the compute nodes. This layer is responsible for most of the I/O operations and saves extreme amounts of data movement between compute and storage resources. We have designed and implemented a system prototype of this architecture - which we call the FusionFS distributed file system - to support metadata-intensive and write-intensive operations, both of which are critical to the I/O performance of scientific applications. FusionFS has been deployed and evaluated on up to 16K compute nodes of an IBM Blue Gene/P supercomputer, showing more than an order of magnitude performance improvement over other popular file systems such as GPFS, PVFS, and HDFS. Dongfang Zhao 0001, Zhao Zhang 0007, Xiaobing Zhou, Tonglin Li, Ke Wang 0012, Dries Kimpe, Philip H. Carns, Robert B. Ross, Ioan Raicu |
IEEE BigData | 1 |
| 2014 | HyCache+: Towards Scalable High-Performance Caching Middleware for Parallel File SystemsabstractThe ever-growing gap between the computation and I/O is one of the fundamental challenges for future computing systems. This computation-I/O gap is even larger for modern large scale high-performance systems due to their state-of-the-art yet decades long architecture: the compute and storage resources form two cliques that are interconnected with shared networking infrastructure. This paper presents a distributed storage middleware, called HyCache+, right on the compute nodes, which allows I/O to effectively leverage the high bi-section bandwidth of the high-speed interconnect of massively parallel high-end computing systems. HyCache+ provides the POSIX interface to end users with the memory-class I/O throughput and latency, and transparently swap the cached data with the existing slow speed but high-capacity networked attached storage. HyCache+ has the potential to achieve both high performance and low cost large capacity, the best of both worlds. To further improve the caching performance from the perspective of the global storage system, we propose a 2-phase mechanism to cache the hot data for parallel applications, called 2-Layer Scheduling (2LS), which minimizes the file size to be transferred between compute nodes and heuristically replaces files in the cache. We deploy HyCache+ on the IBM Blue Gene/P supercomputer, and observe two orders of magnitude faster I/O throughput than the default GPFS parallel file system. Furthermore, the proposed heuristic caching approach shows 29X speedup over the traditional LRU algorithm. Dongfang Zhao 0001, Kan Qiao, Ioan Raicu |
CCGRID | 1 |
| 2013 | Towards high-performance and cost-effective distributed storage systems with information dispersal algorithmsabstractReliability is one of the most fundamental challenges for high performance computing (HPC) and cloud computing. Data replication is the de facto mechanism to achieve high reliability, even though it has been criticized for its high cost and low efficiency. Recent research showed promising results by switching the traditional data replication to a software-based RAID. In order to systematically study the effectiveness of this new method, we built two storage systems from the ground up: a POSIX-compliant distributed file system (FusionFS) and a distributed key-value store (IStore), both supporting information dispersal algorithms (IDA) for data redundancy. FusionFS is crafted to have excellent throughput and scalability for HPC, whereas IStore is architected mainly as a light-weight key-value storage in cloud computing. We evaluated both systems with a large number of parameter combinations. Results show that, for both HPC and cloud computing communities, IDA-based methods with current commodity hardware could outperform data replication in some cases, and would completely surpass data replication with the growing computational capacity through multi/many-core processors (e.g. Intel Xeon Phi, NVIDIA GPU). Dongfang Zhao 0001, Kent Burlingame, Corentin Debains, Pedro Alvarez-Tabio, Ioan Raicu |
CLUSTER | 1 |
| 2013 | Distributed data provenance for large-scale data-intensive computingabstractIt has become increasingly important to capture and understand the origins and derivation of data (its provenance). A key issue in evaluating the feasibility of data provenance is its performance, overheads, and scalability. In this paper, we explore the feasibility of a general metadata storage and management layer for parallel file systems, in which metadata includes both file operations and provenance metadata. We experimentally investigate the design optimality—whether provenance metadata should be loosely-coupled or tightly integrated with a file metadata storage systems. We consider two systems that have applied similar distributed concepts to metadata management, but focusing singularly on kind of metadata: (i) FusionFS, which implements a distributed file metadata management based on distributed hash tables, and (ii) SPADE, which uses a graph database to store audited provenance data and provides distributed module for querying provenance. Our results on a 32-node cluster show that FusionFS+SPADE is a promising prototype with negligible provenance overhead and has promise to scale to petascale and beyond. Furthermore, FusionFS with its own storage layer for provenance capture is able to scale up to 1K nodes on BlueGene/P supercomputer. Dongfang Zhao 0001, Chen Shou, Tanu Malik, Ioan Raicu |
CLUSTER | 1 |
| 2013 | ZHT: A Light-Weight Reliable Persistent Dynamic Scalable Zero-Hop Distributed Hash TableabstractThis paper presents ZHT, a zero-hop distributed hash table, which has been tuned for the requirements of high-end computing systems. ZHT aims to be a building block for future distributed systems, such as parallel and distributed file systems, distributed job management systems, and parallel programming systems. The goals of ZHT are delivering high availability, good fault tolerance, high throughput, and low latencies, at extreme scales of millions of nodes. ZHT has some important properties, such as being light-weight, dynamically allowing nodes join and leave, fault tolerant through replication, persistent, scalable, and supporting unconventional operations such as append (providing lock-free concurrent key/value modifications) in addition to insert/lookup/remove. We have evaluated ZHT's performance under a variety of systems, ranging from a Linux cluster with 512-cores, to an IBM Blue Gene/P supercomputer with 160K-cores. Using micro-benchmarks, we scaled ZHT up to 32K-cores with latencies of only 1.1ms and 18M operations/sec throughput. This work provides three real systems that have integrated with ZHT, and evaluate them at modest scales. 1) ZHT was used in the FusionFS distributed file system to deliver distributed meta-data management at over 60K operations (e.g. file create) per second at 2K-core scales. 2) ZHT was used in the IStore, an information dispersal algorithm enabled distributed object storage system, to manage chunk locations, delivering more than 500 chunks/sec at 32-nodes scales. 3) ZHT was also used as a building block to MATRIX, a distributed job scheduling system, delivering 5000 jobs/sec throughputs at 2K-core scales. We compared ZHT against other distributed hash tables and key/value stores and found it offers superior performance for the features and portability it supports. Tonglin Li, Xiaobing Zhou, Kevin Brandstatter, Dongfang Zhao 0001, Ke Wang 0012, Anupam Rajendran, Zhao Zhang 0007, Ioan Raicu |
IPDPS | 4 |
| 2009 | Incremental Isometric Embedding of High-Dimensional Data Using Connected Neighborhood GraphsabstractMost nonlinear data embedding methods use bottom-up approaches for capturing the underlying structure of data distributed on a manifold in high dimensional space. These methods often share the first step which defines neighbor points of every data point by building a connected neighborhood graph so that all data points can be embedded to a single coordinate system. These methods are required to work incrementally for dimensionality reduction in many applications. Because input data stream may be under-sampled or skewed from time to time, building connected neighborhood graph is crucial to the success of incremental data embedding using these methods. This paper presents algorithms for updating $k$-edge-connected and $k$-connected neighborhood graphs after a new data point is added or an old data point is deleted. It further utilizes a simple algorithm for updating all-pair shortest distances on the neighborhood graph. Together with incremental classical multidimensional scaling using iterative subspace approximation, this paper devises an incremental version of Isomap with enhancements to deal with under-sampled or unevenly distributed data. Experiments on both synthetic and real-world data sets show that the algorithm is efficient and maintains low dimensional configurations of high dimensional data under various data distributions. Dongfang Zhao 0001 |
IEEE Trans. Pattern Anal. Mach. Intell. | 1 |
| 2008 | Solving SQL Constraints by Incremental Translation to SAT
Robin Lohfert, James J. Lu, Dongfang Zhao 0001 |
IEA/AIE | 3 |