EDBT 2026 Demo / reviewers in the wild / expert
Zhao Zhang 0007
dblp:87/6853-7
· DBLP profile ↗
37ranked-venue papers
9as first author
14since 2021 · last 2026
0000-0001-5921-0035ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 30 · 6 first-author · 13 since 2021Applied, interdisciplinary, general and emerging computing · 5 · 3 first-author · 1 since 2021Databases, data management, data science and information retrieval · 4 · 1 first-authorArtificial intelligence and machine learning · 2 · 1 first-authorSoftware engineering, systems software and programming languages · 2 · 1 first-author · 1 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Ciri: Scaling-aware job scheduling for small language models on over-committed GPU clustersabstractDespite the remarkable breakthroughs in large language models (LLMs), researchers are also designing models with modest size (hundreds of millions to billions of parameters) for domain science and agentic AI. Training and fine-tuning modest models on over-committed GPU clusters can experience long turnarounds due to the high computational demand and extensive queue wait times. Enabled by gradient accumulation, such a workload can leverage any number of GPUs as long as the global mini-batch size is fixed. Training with more GPUs can reduce execution time, but leads to longer queue wait times and higher node hour consumption due to nonlinear scaling efficiency. There is a clear tradeoff between turnaround time and the node hour consumption, which is agnostic to existing schedulers. In this paper, we present Ciri, a reinforcement learning-enabled and scaling-aware user-level job scheduler, which incorporates workload scaling curves and machine states into scheduling decisions. Ciri dynamically selects the GPU count to trade off turnaround time for node hour consumption. We evaluate Ciri with three GPT models: GPT-345M, GPT-1.5B, and GPT-2.0B with job traces of the TACC Lonestar6 and Frontera GPU clusters. The results show that Ciri can effectively reduce the node hour consumption by 0.9-1.8% for every 1% increase in turnaround across the two machines. Shuyuan Fan, Zhao Zhang 0007 |
Future Gener. Comput. Syst. | 3 |
| 2026 | Quantifying Performance Variability in GPU ClustersabstractModern supercomputers are equipped with massive amounts of Graphics Processing Units (GPUs) to meet the growing demands of scientific computing and machine learning. However, GPUs, even of the same type, exhibit performance variability, leading to resource under-utilization and prolonged execution times. In this work, we perform a characterization study on the performance variability of NVIDIA A100 GPUs and GH200 superchips using the GEMM and STREAM micro-benchmarks and seven real-world applications across two supercomputing systems: TACC Vista (GH200) and NERSC Perlmutter (A100). Our results reveal GEMM performance variability ranging from 0.1% to 8.8%, though outliers can push deviations significantly higher. For scientific applications, the variability ranges from 1.0% to 12.2% on Vista. The variability of single-GPU GPT 4.8B and Llama 3B training is 10.4% and 8.9%, respectively. We further extend this study by training a GPT 19B model using eight GH200s in the 3D parallel configuration, observing up to 3.6% throughput variability and 8.5% slowdown in the worst case. By comparing the variability across GPU types, compute paths, data precisions, and applications, we observe that applications that use Tensor Cores and FP64 on CUDA cores show higher variability than other applications. Supercomputer users can leverage the observed variability to diagnose performance issues and enhance execution efficiency. Computing centers can design variability-aware scheduling approaches to achieve higher machine utilization without sacrificing individual application performance. Michael Mogilevsky, Hazem Zaky, Mingkai Zheng, Lishan Yang 0001, Zhao Zhang 0007 |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2025 | Diamond: Harnessing GPU Resources for Scientific Deep LearningabstractModern research computing cyberinfrastructure, such as ACCESS-CI and NAIRR Pilot, offers GPU resources across geographically distributed clusters to accommodate the increasing needs of scientific deep learning (DL) workloads. Even for high-performance computing (HPC) experts, configuring environments and managing DL workloads across supercomputers remain significant barriers. To address these obstacles, we present Diamond, an open-source platform to simplify and streamline the DL lifecycle on HPC. Diamond provides an intuitive graphical interface that abstracts system-level complexity, enabling users to develop, debug, and deploy DL models with minimal overhead. We identify several challenges in building such a platform, including portability, security, and usability, and propose effective architectural solutions to each. Notably, Diamond enables users to share and reuse DL workload environments across systems and collaborators, reducing redundant setup efforts. Experimental results demonstrate that Diamond reduces the time to first successful deployment by an average of 68%, compared to manual configuration with command lines. The Diamond service is available at https://diamondhpc.ai. Haotian Xie, Rohan Marwaha, Minu Mathew, Song Bian 0002, Gengcong Yang, Minghao Yan, Yadu N. Babuji, Owen Price, Yinzhi Wang, Volodymyr V. Kindratenko, Shivaram Venkataraman, Kyle Chard, Ian T. Foster, Zhao Zhang 0007 |
eScience | 14 |
| 2025 | Efficient Fine-Grained Gpu Performance Modeling for Distributed Deep Learning of LlmabstractTraining Large Language Models(LLMs) is one of the most compute-intensive tasks in high-performance computing. Predicting end-to-end training time for multi-billion parameter models distributed across hundreds of GPUs remains challenging due to complex interactions between transformer components, parallelism strategies(data, model, pipeline, tensor), and multi-tier communication. Learned models require costly sampling, while analytical models often struggle with real-world network and hardware complexities. We address this by decomposing LLMs into core computational primitives and modeling them with: (1) operator-level decomposition for fine-grained analysis; (2) lightweight sampling based hardware-aware prediction models for key operations; (3) an end-to-end prediction system integrating these components across complex parallelization strategies. Crucially, our methodology has been validated on two large-scale HPC systems. Our framework achieves low average prediction errors-4.98% on Perlmutter(A100) and 9.38% on Vista(GH200)-for models up to 20B parameters across 128 GPUs. Importantly, it runs entirely on CPUs, enabling rapid iteration over hardware configurations and training strategies without costly on-cluster experimentation. Biyao Zhang, Mingkai Zheng, Debargha Ganguly, Xuecen Zhang, Vikash Singh, Vipin Chaudhary, Zhao Zhang 0007 |
HiPC | 7 |
| 2025 | FT2: First-Token-Inspired Online Fault Tolerance on Critical Layers for Generative Large Language ModelsabstractGenerative Large Language Models (LLMs) are deployed on large-scale computing systems, where such tasks unavoidably suffer from soft errors, leading to quality degradation of content generated by LLMs. Enhancing LLM resilience is particularly challenging because of its complicated model architecture and tremendous size. State-of-the-art protections have limitations such as high overhead and incomplete coverage, and often require offline profiling. Cherish Mulpuru, Roberto Gioiosa, Zhao Zhang 0007, Bo Fang 0002, Lishan Yang 0001 |
HPDC | 5 |
| 2025 | COMPSO: Optimizing Gradient Compression for Distributed Training with Second-Order OptimizersabstractSecond-order optimization methods have been developed to enhance convergence and generalization in deep neural network (DNN) training compared to first-order methods like Stochastic Gradient Descent (SGD). However, these methods face challenges in distributed settings due to high communication overhead. Gradient compression, a technique commonly used to accelerate communication for first-order approaches, often results in low communication reduction ratios, decreased model accuracy, and/or high compression overhead when applied to second-order methods. To address these limitations, we introduce a novel gradient compression method for second-order optimizers called COMPSO. This method effectively reduces communication costs while preserving the advantages of second-order optimization. COMPSO employs stochastic rounding to maintain accuracy and filters out minor gradients to improve compression ratios. Additionally, we develop GPU optimizations to minimize compression overhead and performance modeling to ensure end-to-end performance gains across various systems. Evaluation of COMPSO on different DNN models shows that it achieves a compression ratio of 22.1×, reduces communication time by 14.2×, and improves overall performance by 1.9×, all without any drop in model accuracy. Baixi Sun, Weijin Liu, J. Gregory Pauloski, Jiannan Tian, Jinda Jia, Daoce Wang, Boyuan Zhang 0002, Mingkai Zheng, Sheng Di, Sian Jin, Zhao Zhang 0007, Xiaodong Yu 0001, Kamil Iskra, Pete Beckman, Guangming Tan, Dingwen Tao |
PPoPP | 11 |
| 2025 | Demystifying the Resilience of Large Language Model Inference: An End-to-End PerspectiveabstractDeep neural networks are known to be resilient to random bitwise faults in their parameters. However, this resilience has primarily been established through studies of classification models. The extent to which this claim holds for large-language models remains under-explored. In this work, we conduct an extensive measurement study on the impact of random bitwise faults in commercial-scale language model inference. We first expose that these language models are not truly resilient to random bit-flips. While aggregate metrics such as accuracy may suggest resilience, an in-depth inspection of the generated outputs shows significant degradation in text quality. Our analysis also shows that tasks requiring more complex reasoning suffer more from performance and quality degradation. Moreover, we extend our resilience analysis to models with augmented reasoning capabilities, such as Chain-of-Thought or Mixture of Experts architectures. Zachary Coalson, Shiyang Chen 0004, Hang Liu 0001, Zhao Zhang 0007, Sanghyun Hong 0001, Bo Fang 0002, Lishan Yang 0001 |
SC | 5 |
| 2024 | Dual Channel Dual Staging: Hierarchical and Portable Staging for GPU-Based In-Situ WorkflowabstractIn-situ workflows have emerged as an attractive approach for addressing data movement challenges at very large scales. Since GPU-based architectures dominate the HPC landscapes, porting these in-situ workflows, and, specifically, the inter-application data exchange, to GPU-based systems can be challenging. Technologies such as GPUDirect RDMA (GDR), which is typically used for I/O in GPU applications as an optimization that circumvents the CPU overhead, can be leveraged to support bulk data exchanges between GPU applications. However, current GDR design often lacks performance portability across HPC clusters built with different hardware configurations. Furthermore, the local CPU may also be effectively used as an auxiliary communication mechanism to offload data exchanges. In this paper, we present a dual channel dual staging approach for efficient, scalable, and performance-portable inter-application data exchange for in-situ workflows. This approach exploits the data access pattern within in-situ workflows along with the inherent execution asynchrony to accelerate data exchanges and, at the same time, improve performance portability. Specifically, the dual channel dual staging method leverages both the local CPU and the remote data staging server to build a hierarchical joint staging area and uses this staging area to transform blocking inter-application bulk data exchanges into best-effort local data movements between GPU and CPU. The dual channel dual staging is implemented as a portability extension of the Dataspaces-GPU staging framework. We present an experimental evaluation of its performance, portability, and scalability using this implementation on three leadership GPU clusters. The evaluation results demonstrate that the dual channel dual staging method saves up to 75% in data-exchange time compared to host-based, GDR, and alternate portable designs, while maintaining scalability (up to 512 GPUs) and performance portability across the three platforms. Bo Zhang 0120, Philip E. Davis, Zhao Zhang 0007, Keita Teranishi, Manish Parashar |
HiPC | 3 |
| 2023 | Optimizing Data Movement for GPU-Based In-Situ Workflow Using GPUDirect RDMA
Bo Zhang 0120, Philip E. Davis, Nicolas M. Morales, Zhao Zhang 0007, Keita Teranishi, Manish Parashar |
Euro-Par | 4 |
| 2023 | Mirage: Towards Low-interruption Services on Batch GPU Clusters with Reinforcement LearningabstractAccommodating long-running deep learning (DL) training and inference jobs is challenging on GPU clusters that use traditional batch schedulers, such as Slurm. Given fixed wall clock time limits, DL researchers usually need to run a sequence of batch jobs and experience long interruptions on overloaded machines. Such interruptions significantly lower the research productivity and QoS for services that are deployed in production. To mitigate the issues from interruption, we propose the design of a proactive provisioner and investigate a set of statistical learning and reinforcement learning (RL) techniques, including random forest, xgboost, Deep Q-Network, and policy gradient. Using production job traces from three GPU clusters, we train each model using a subset of the trace and then evaluate their generality using the remaining validation subset. We introduce Mirage, a Slurm-compatible resource provisioner that integrates the candidate ML methods. Our experiments show that the Mirage can reduce interruption by 17--100% and safeguard 23%-76% of jobs with zero interruption across varying load levels on the three clusters. Qiyang Ding, Shreyas Kudari, Shivaram Venkataraman, Zhao Zhang 0007 |
SC | 5 |
| 2023 | Fine-grained Policy-driven I/O Sharing for Burst BuffersabstractA burst buffer is a common method to bridge the performance gap between the I/O needs of modern supercomputing applications and the performance of the shared file system on large-scale supercomputers. However, existing I/O sharing methods require resource isolation, offline profiling, or repeated execution that significantly limit the utilization and applicability of these systems. Here we present ThemisIO, a policy-driven I/O sharing framework for a remote-shared burst buffer: a dedicated group of I/O nodes, each with a local storage device. ThemisIO preserves high utilization by implementing opportunity fairness so that it can reallocate unused I/O resources to other applications. ThemisIO accurately and efficiently allocates I/O cycles among applications, purely based on real-time I/O behavior without requiring user-supplied information or offline-profiled application characteristics. ThemisIO supports a variety of fair sharing policies, such as user-fair, size-fair, as well as composite policies, e.g., group-then-user-fair. All these features are enabled by its statistical token design. ThemisIO can alter the execution order of incoming I/O requests based on assigned tokens to precisely balance I/O cycles between applications via time slicing, thereby enforcing processing isolation. Experiments using I/O benchmarks show that ThemisIO sustains 13.5--13.7% higher I/O throughput and 19.5--40.4% lower performance variation than existing algorithms. For real applications, ThemisIO significantly reduces the slowdown by 59.1--99.8% caused by I/O interference. Ed Karrels, Lei Huang 0019, Yuhong Kan, Ishank Arora, Yinzhi Wang, Daniel S. Katz, William Gropp, Zhao Zhang 0007 |
SC | 8 |
| 2022 | Deep Neural Network Training With Distributed K-FACabstractScaling deep neural network training to more processors and larger batch sizes is key to reducing end-to-end training time; yet, maintaining comparable convergence and hardware utilization at larger scales is challenging. Increases in training scales have enabled natural gradient optimization methods as a reasonable alternative to stochastic gradient descent and variants thereof. Kronecker-factored Approximate Curvature (K-FAC), a natural gradient method, preconditions gradients with an efficient approximation of the Fisher Information Matrix to improve per-iteration progress when optimizing an objective function. Here we propose a scalable K-FAC algorithm and investigate K-FAC’s applicability in large-scale deep neural network training. Specifically, we explore layer-wise distribution strategies, inverse-free second-order gradient evaluation, and dynamic K-FAC update decoupling, with the goal of preserving convergence while minimizing training time. We evaluate the convergence and scaling properties of our K-FAC gradient preconditioner, for image classification, object detection, and language modeling applications. In all applications, our implementation converges to baseline performance targets in 9–25% less time than the standard first-order optimizers on GPU clusters across a variety of scales. J. Gregory Pauloski, Lei Huang 0019, Weijia Xu, Kyle Chard, Ian T. Foster, Zhao Zhang 0007 |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2021 | Characterizing Impacts of Storage Faults on HPC Applications: A Methodology and InsightsabstractIn recent years, the increasing complexity in scientific simulations and emerging demands for training heavy artificial intelligence models require massive and fast data accesses, which urges high-performance computing (HPC) platforms to equip with more advanced storage infrastructures such as solid-state disks (SSDs). While SSDs offer high-performance I/O, the reliability challenges faced by the HPC applications under the SSD-related failures remains unclear, in particular for failures resulting in data corruptions. The goal of this paper is to understand the impact of SSD-related faults on the behaviors of complex HPC applications. To this end, we propose FFIS, a FUSE-based fault injection framework that systematically introduces storage faults into the application layer to model the errors originated from SSDs. FFIS is able to plant different I/O related faults into the data returned from underlying file systems, which enables the investigation on the error resilience characteristics of the scientific file format. We demonstrate the use of FFIS with three representative real HPC applications, showing how each application reacts to the data corruptions, and provide insights on the error resilience of the widely adopted HDF5 file format for the HPC applications. Bo Fang 0002, Daoce Wang, Sian Jin, Quincey Koziol, Zhao Zhang 0007, Qiang Guan, Surendra Byna, Sriram Krishnamoorthy, Dingwen Tao |
CLUSTER | 5 |
| 2021 | KAISA: an adaptive second-order optimizer framework for deep neural networksabstractKronecker-factored Approximate Curvature (K-FAC) has recently been shown to converge faster in deep neural network (DNN) training than stochastic gradient descent (SGD); however, K-FAC's larger memory footprint hinders its applicability to large models. We present KAISA, a K-FAC-enabled, Adaptable, Improved, and ScAlable second-order optimizer framework that adapts the memory footprint, communication, and computation given specific models and hardware to improve performance and increase scalability. We quantify the tradeoffs between memory and communication cost and evaluate KAISA on large models, including ResNet-50, Mask R-CNN, U-Net, and BERT, on up to 128 NVIDIA A100 GPUs. Compared to the original optimizers, KAISA converges 18.1--36.3% faster across applications with the same global batch size. Under a fixed memory budget, KAISA converges 32.5% and 41.6% faster in ResNet-50 and BERT-Large, respectively. KAISA can balance memory and communication to achieve scaling efficiency equal to or better than the baseline optimizers. J. Gregory Pauloski, Lei Huang 0019, Shivaram Venkataraman, Kyle Chard, Ian T. Foster, Zhao Zhang 0007 |
SC | 7 |
| 2020 | Efficient I/O for Neural Network Training with Compressed DataabstractFanStore is a shared object store that enables efficient and scalable neural network training on supercomputers. By providing a global cache layer on node-local burst buffers using a compressed representation, it significantly enhances the processing capability of deep learning (DL) applications on existing hardware. In addition, FanStore allows POSIX-compliant file access to the compressed data in user space. We investigate the tradeoff between runtime overhead and data compression ratio using real-world datasets and applications, and propose a compressor selection algorithm to maximize storage capacity given performance constraints. We consider both asynchronous (i.e., with prefetching) and synchronous I/O strategies, and propose mechanisms for selecting compressors for both approaches. Using FanStore, the same storage hardware can host 2–13× more data for example applications without significant runtime overhead. Empirically, our experiments show that FanStore scales to 512 compute nodes with near linear performance scalability. Zhao Zhang 0007, Lei Huang 0019, J. Gregory Pauloski, Ian T. Foster |
IPDPS | 1 |
| 2020 | Convolutional neural network training with distributed K-FACabstractTraining neural networks with many processors can reduce time-to-solution; however, it is challenging to maintain convergence and efficiency at large scales. The Kroneckerfactored Approximate Curvature (K-FAC) was recently proposed as an approximation of the Fisher Information Matrix that can be used in natural gradient optimizers. We investigate here a scalable K-FAC design and its applicability in convolutional neural network (CNN) training at scale. We study optimization techniques such as layer-wise distribution strategies, inverse-free second-order gradient evaluation, and dynamic K-FAC update decoupling to reduce training time while preserving convergence. We use residual neural networks (ResNet) applied to the CIFAR10 and ImageNet-1k datasets to evaluate the correctness and scalability of our K-FAC gradient preconditioner. With ResNet-50 on the ImageNet-1k dataset, our distributed K-FAC implementation converges to the 75.9% MLPerf baseline in 18-25% less time than does the classic stochastic gradient descent (SGD) optimizer across scales on a GPU cluster. J. Gregory Pauloski, Zhao Zhang 0007, Lei Huang 0019, Weijia Xu, Ian T. Foster |
SC | 2 |
| 2020 | Kira: Processing Astronomy Imagery Using Big Data TechnologyabstractScientific analyses commonly compose multiple single-process programs into a dataflow. An end-to-end dataflow of single-process programs is known as a many-task application. Typically, HPC tools are used to parallelize these analyses. In this work, we investigate an alternate approach that uses Apache Spark-a modern platform for data intensive computing-to parallelize many-task applications. We implement Kira, a flexible and distributed astronomy image processing toolkit, and its Source Extractor (Kira SE) application. Using Kira SE as a case study, we examine the programming flexibility, dataflow richness, scheduling capacity and performance of Apache Spark running on the Amazon EC2 cloud. By exploiting data locality, Kira SE achieves a 4.1× speedup over an equivalent C program when analyzing a 1TB dataset using 512 cores on the Amazon EC2 cloud. Furthermore, Kira SE on the Amazon EC2 cloud achieves a 1.8× speedup over the C program on the NERSC Edison supercomputer. A 128-core Amazon EC2 cloud deployment of Kira SE using Spark Streaming can achieve a second-scale latency with a sustained throughput of 800 MB/s. Our experience with Kira demonstrates that data intensive computing platforms like Apache Spark are a performant alternative for many-task scientific applications. Zhao Zhang 0007, Kyle Barbary, Frank A. Nothaft, Evan Randall Sparks, Oliver Zahn, Michael J. Franklin, David A. Patterson 0001, Saul Perlmutter |
IEEE Trans. Big Data | 1 |
| 2019 | Quantifying the Impact of Memory Errors in Deep LearningabstractThe use of deep learning (DL) on HPC resources has become common as scientists explore and exploit DL methods to solve domain problems. On the other hand, in the coming exascale computing era, a high error rate is expected to be problematic for most HPC applications. The impact of errors on DL applications, especially DL training, remains unclear given their stochastic nature. In this paper, we focus on understanding DL training applications on HPC in the presence of silent data corruption. Specifically, we design and perform a quantification study with three representative applications by manually injecting silent data corruption errors (SDCs) across the design space and compare training results with the error-free baseline. The results show only 0.61-1.76% of SDCs cause training failures, and taking the SDC rate in modern hardware into account, the actual chance of a failure is one in thousands to millions of executions. With this quantitatively measured impact, computing centers can make rational design decisions based on their application portfolio, the acceptable failure rate, and financial constraints; for example, they might determine their confidence in the correctness of training results performed on processors without error correction code (ECC) RAM. We also discover that over 75-90% of the SDCs that cause catastrophic errors can be easily detected by a training loss in the next iteration. Thus we propose this error-aware software solution to correct catastrophic errors, as it has significantly lower time and space overhead compared to algorithm-based fault-tolerance (ABFT) and ECC. Zhao Zhang 0007, Lei Huang 0019, Ruizhu Huang, Weijia Xu, Daniel S. Katz |
CLUSTER | 1 |
| 2019 | Fast Deep Neural Network Training on Distributed Systems and Cloud TPUsabstractSince its creation, the ImageNet-1k benchmark set has played a significant role as a benchmark for ascertaining the accuracy of different deep neural net (DNN) models on the image classification problem. Moreover, in recent years it has also served as the principal benchmark for assessing different approaches to DNN training. Finishing a 90-epoch ImageNet-1k training with ResNet-50 on a NVIDIA M40 GPU takes 14 days. This training requires 1018 single precision operations in total. On the other hand, the world's current fastest supercomputer can finish 3 x 1017 single precision operations per second (according to the Nov 2018 Top 500 results). If we can make full use of the computing capability of the fastest supercomputer, we should be able to finish the training in several seconds. Over the last two years, researchers have focused on closing this significant performance gap through scaling DNN training to larger numbers of processors. Most successful approaches to scaling ImageNet training have used the synchronous minibatch stochastic gradient descent (SGD). However, to scale synchronous SGD one must also increase the batch size used in each iteration. Thus, for many researchers, the focus on scaling DNN training has translated into a focus on developing training algorithms that enable increasing the batch size in data-parallel synchronous SGD without losing accuracy over a fixed number of epochs. In this paper, we investigate supercomputers' capability of speeding up DNN training. Our approach is to use a large batch size, powered by the Layer-wise Adaptive Rate Scaling (LARS) algorithm, for efficient usage of massive computing resources. Our approach is generic, as we empirically evaluate the effectiveness on five neural networks: AlexNet, AlexNet-BN, GNMT, ResNet-50, and ResNet-50-v2 trained with large datasets while preserving the state-of-the-art test accuracy. Compared to the baseline of a previous study from Goyal et al. [1], our approach shows higher test accuracy on batch sizes that are larger than 16K.When we use the same baseline, our results are better than Goyal et al. for all the batch sizes (Fig. 20). Using 2,048 Intel Xeon Platinum 8160 processors, we reduce the 100-epoch AlexNet training time from hours to 11 minutes. With 2,048 Intel Xeon Phi 7250 Processors, we reduce the 90-epoch ResNet-50 training time from hours to 20 minutes. Our implementation is open source and has been released in the Intel distribution of Caffe, Facebook's PyTorch, and Google's TensorFlow. The difference between this paper and the conference-version of our work [2] includes: (1) we implement our approach on Google's cloud Tensor Processing Unit (TPU) platform, which verifies our previous success on CPUs and GPUs. (2) we scale the batch size of ResNet-50-v2 to 32K and achieve 76.3 percent accuracy, which is better than the 75.3 percent accuracy achieved in our conference paper. (3) we apply our approach to Google's Neural Machine Translation (GNMT) application, which helps us to achieves 4x speedup on the cloud TPUs. Yang You 0001, Zhao Zhang 0007, Cho-Jui Hsieh, James Demmel, Kurt Keutzer |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2018 | BeeFlow: A Workflow Management System for In Situ Processing across HPC and Cloud SystemsabstractIn this paper, we propose BeeFlow - an in situ analysis enabled workflow management system across multiple platforms using Docker containers. BeeFlow can support both traditional workflows as well as workflows with in situ analysis. BeeFlow leverages Docker containers to provide a portable, flexible, and reproducible workflow management system across HPC and cloud platforms. We showcase how current in situ visualization workflows can apply BeeFlow with DOE production codes VPIC and Flecsale. Jieyang Chen, Qiang Guan, Zhao Zhang 0007, Xin Liang 0001, Louis James Vernon, Allen McPherson, Li-Ta Lo, Patricia Grubel, Tim Randles, Zizhong Chen, James P. Ahrens |
ICDCS | 3 |
| 2018 | ImageNet Training in MinutesabstractIn this paper, we investigate large scale computers' capability of speeding up deep neural networks (DNN) training. Our approach is to use large batch size, powered by the Layer-wise Adaptive Rate Scaling (LARS) algorithm, for efficient usage of massive computing resources. Our approach is generic, as we empirically evaluate the effectiveness on two neural networks: AlexNet and ResNet-50 trained with the ImageNet-1k dataset while preserving the state-of-the-art test accuracy. Compared to the baseline of a previous study from a group of researchers at Facebook, our approach shows higher test accuracy on batch sizes that are larger than 16K. Using 2,048 Intel Xeon Platinum 8160 processors, we reduce the 100-epoch AlexNet training time from hours to 11 minutes. With 2,048 Intel Xeon Phi 7250 Processors, we reduce the 90-epoch ResNet-50 training time from hours to 20 minutes. Our implementation is open source and has been released in the Intel distribution of Caffe v1.0.7. Yang You 0001, Zhao Zhang 0007, Cho-Jui Hsieh, James Demmel, Kurt Keutzer |
ICPP | 2 |
| 2017 | Diagnosing Machine Learning Pipelines with Fine-grained LineageabstractWe present the Hippo system to enable the diagnosis of distributed machine learning (ML) pipelines by leveraging fine-grained data lineage. Hippo exposes a concise yet powerful API, derived from primitive lineage types, to capture fine-grained data lineage for each data transformation. It records the input datasets, the output datasets and the cell-level mapping between them. It also collects sufficient information that is needed to reproduce the computation. Hippo efficiently enables common ML diagnosis operations such as code debugging, result analysis, data anomaly removal, and computation replay. By exploiting the metadata separation and high-order function encoding strategies, we observe an O(10^3)x total improvement in lineage storage efficiency vs. the baseline of cell-wise mapping recording while maintaining the lineage integrity. Hippo can answer the real use case lineage queries within a few seconds, which is low enough to enable interactive diagnosis of ML pipelines. Zhao Zhang 0007, Evan Randall Sparks, Michael J. Franklin |
HPDC | 1 |
| 2016 | Integrating Abstractions to Enhance the Execution of Distributed ApplicationsabstractOne of the factors that limits the scale, performance, and sophistication of distributed applications is the difficulty of concurrently executing them on multiple distributed computing resources. In part, this is due to a poor understanding of the general properties and performance of the coupling between applications and dynamic resources. This paper addresses this issue by integrating abstractions representing distributed applications, resources, and execution processes into a pilot-based middleware. The middleware provides a platform that can specify distributed applications, execute them on multiple resource and for different configurations, and is instrumented to support investigative analysis. We analyzed the execution of distributed applications using experiments that measure the benefits of using multiple resources, the late-binding of scheduling decisions, and the use of backfill scheduling. Matteo Turilli, Zhao Zhang 0007, André Merzky, Michael Wilde, Jon B. Weissman, Daniel S. Katz, Shantenu Jha |
IPDPS | 3 |
| 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. | 6 |
| 2016 | Application skeletons: Construction and use in eScience
Daniel S. Katz, André Merzky, Zhao Zhang 0007, Shantenu Jha |
Future Gener. Comput. Syst. | 3 |
| 2015 | Scientific computing meets big data technology: An astronomy use caseabstractScientific analyses commonly compose multiple single-process programs into a dataflow. An end-to-end dataflow of single-process programs is known as a many-task application. Typically, tools from the HPC software stack are used to parallelize these analyses. In this work, we investigate an alternate approach that uses Apache Spark - a modern big data platform - to parallelize many-task applications. We present Kira, a flexible and distributed astronomy image processing toolkit using Apache Spark. We then use the Kira toolkit to implement a Source Extractor application for astronomy images, called Kira SE. With Kira SE as the use case, we study the programming flexibility, dataflow richness, scheduling capacity and performance of Apache Spark running on the EC2 cloud. By exploiting data locality, Kira SE achieves a 3.7 χ speedup over an equivalent C program when analyzing a 1TB dataset using 512 cores on the Amazon EC2 cloud. Furthermore, we show that by leveraging software originally designed for big data infrastructure, Kira SE achieves competitive performance to the C implementation running on the NERSC Edison supercomputer. Our experience with Kira indicates that emerging Big Data platforms such as Apache Spark are a performant alternative for many-task scientific applications. Zhao Zhang 0007, Kyle Barbary, Frank A. Nothaft, Evan Randall Sparks, Oliver Zahn, Michael J. Franklin, David A. Patterson 0001, Saul Perlmutter |
IEEE BigData | 1 |
| 2015 | The Missing Piece in Complex Analytics: Low Latency, Scalable Model Management and Serving with Velox
Daniel Crankshaw, Peter Bailis, Joseph Gonzalez 0001, Haoyuan Li 0001, Zhao Zhang 0007, Michael J. Franklin, Ali Ghodsi 0002, Michael I. Jordan |
CIDR | 5 |
| 2015 | Rethinking Data-Intensive Science Using Scalable Analytics Systemsabstract"Next generation" data acquisition technologies are allowing scientists to collect exponentially more data at a lower cost. These trends are broadly impacting many scientific fields, including genomics, astronomy, and neuroscience. We can attack the problem caused by exponential data growth by applying horizontally scalable techniques from current analytics systems to accelerate scientific processing pipelines. Frank A. Nothaft, Matt Massie, Timothy Danford, Zhao Zhang 0007, Uri Laserson, Carl Yeksigian, Jey Kottalam, Arun Ahuja, Jeff Hammerbacher, Michael D. Linderman, Michael J. Franklin, Anthony D. Joseph, David A. Patterson 0001 |
SIGMOD Conference | 4 |
| 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 | 2 |
| 2014 | Using Application Skeletons to Improve eScience InfrastructureabstractComputer scientists who work on tools and systems to support eScience (a variety of parallel and distributed) applications usually use actual applications to prove that their systems will benefit science and engineering (e.g., improve application performance). Accessing and building the applications and necessary data sets can be difficult because of policy or technical issues, and it can be difficult to modify the characteristics of the applications to understand corner cases in the system design. In this paper, we present the Application Skeleton, a simple yet powerful tool to build synthetic applications that represent real applications, with runtime and I/O close to those of the real applications. This allows computer scientists to focus on the system they are building, they can work with the simpler skeleton applications and be sure that their work will also be applicable to the real applications. In addition, skeleton applications support simple reproducible system experiments since they are represented by a compact set of parameters. Our Application Skeleton tool (available as open source at https://github.com/applicationskeleton/Skeleton) currently can create easy-to-access, easy-to-build, and easy-to-run bag-of-task, (iterative) map-reduce, and (iterative) multistage workflow applications. The tasks can be serial or parallel or a mix of both. We select three representative applications (Montage, BLAST, CyberShake Postprocessing), then describe and generate skeleton applications for each. We show that the skeleton applications have identical (or close) performance to that of the real applications. We then show examples of using skeleton applications to verify system optimizations such as data caching, I/O tuning, and task scheduling, as well as the system resilience mechanism, in some cases modifying the skeleton applications to emphasize some characteristic, and thus show that using skeleton applications simplifies the process of designing, implementing, and testing these optimizations. Zhao Zhang 0007, Daniel S. Katz |
eScience | 1 |
| 2014 | Special issue on eScience infrastructure and applications
Daniel S. Katz, Zhao Zhang 0007 |
Future Gener. Comput. Syst. | 2 |
| 2013 | MTC envelope: defining the capability of large scale computers in the context of parallel scripting applications
Zhao Zhang 0007, Daniel S. Katz, Michael Wilde, Justin M. Wozniak, Ian T. Foster |
HPDC | 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 | 7 |
| 2013 | Parallelizing the execution of sequential scriptsabstractScripting is often used in science to create applications via the composition of existing programs. Parallel scripting systems allow the creation of such applications, but each system introduces the need to adopt a somewhat specialized programming model. We present an alternative scripting approach, AMFS Shell, that lets programmers express parallel scripting applications via minor extensions to existing sequential scripting languages, such as Bash, and then execute them in-memory on large-scale computers. We define a small set of commands between the scripts and a parallel scripting runtime system, so that programmers can compose their scripts in a familiar scripting language. The underlying AMFS implements both collective (fast file movement) and functional (transformation based on content) file management. Tasks are handled by AMFS's built-in execution engine. AMFS Shell is expressive enough for a wide range of applications, and the framework can run such applications efficiently on large-scale computers. Zhao Zhang 0007, Daniel S. Katz, Timothy G. Armstrong, Justin M. Wozniak, Ian T. Foster |
SC | 1 |
| 2012 | A Workflow-Aware Storage System: An Opportunity StudyabstractThis paper evaluates the potential gains a workflow-aware storage system can bring. Two observations make us believe such storage system is crucial to efficiently support workflow-based applications: First, workflows generate irregular and application-dependent data access patterns. These patterns render existing storage systems unable to harness all optimization opportunities as this often requires conflicting optimization options or even conflicting design decision at the level of the storage system. Second, when scheduling, workflow runtime engines make suboptimal decisions as they lack detailed data location information. This paper discusses the feasibility, and evaluates the potential performance benefits brought by, building a workflow-aware storage system that supports per-file access optimizations and exposes data location. To this end, this paper presents approaches to determine the application-specific data access patterns, and evaluates experimentally the performance gains of a workflow-aware storage approach. Our evaluation using synthetic benchmarks shows that a workflow-aware storage system can bring significant performance gains: up to 7× performance gain compared to the distributed storage system - MosaStore and up to 16× compared to a central, well provisioned, NFS server. Emalayan Vairavanathan, Samer Al-Kiswany, Lauro Beltrão Costa, Zhao Zhang 0007, Daniel S. Katz, Michael Wilde, Matei Ripeanu |
CCGRID | 4 |
| 2012 | Design and analysis of data management in scalable parallel scriptingabstractWe seek to enable efficient large-scale parallel execution of applications in which a shared filesystem abstraction is used to couple many tasks. Such parallel scripting (many-task computing, MTC) applications suffer poor performance and utilization on large parallel computers because of the volume of filesystem I/O and a lack of appropriate optimizations in the shared filesystem. Thus, we design and implement a scalable MTC data management system that uses aggregated compute node local storage for more efficient data movement strategies. We co-design the data management system with the data-aware scheduler to enable dataflow pattern identification and automatic optimization. The framework reduces the time to solution of parallel stages of an astronomy data analysis application, Montage, by 83.2% on 512 cores; decreases the time to solution of a seismology application, CyberShake, by 7.9% on 2,048 cores; and delivers BLAST performance better than mpiBLAST at various scales up to 32,768 cores, while preserving the flexibility of the original BLAST application. Zhao Zhang 0007, Daniel S. Katz, Justin M. Wozniak, Allan Espinosa, Ian T. Foster |
SC | 1 |
| 2008 | Toward loosely coupled programming on petascale systemsabstractWe have extended the Falkon lightweight task execution framework to make loosely coupled programming on petascale systems a practical and useful programming model. This work studies and measures the performance factors involved in applying this approach to enable the use of petascale systems by a broader user community, and with greater ease. Our work enables the execution of highly parallel computations composed of loosely coupled serial jobs with no modifications to the respective applications. This approach allows a new-and potentially far larger-class of applications to leverage petascale systems, such as the IBM Blue Gene/P supercomputer. We present the challenges of I/O performance encountered in making this model practical, and show results using both microbenchmarks and real applications from two domains: economic energy modeling and molecular dynamics. Our benchmarks show that we can scale up to 160 K processor-cores with high efficiency, and can achieve sustained execution rates of thousands of tasks per second. Ioan Raicu, Zhao Zhang 0007, Michael Wilde, Ian T. Foster, Pete Beckman, Kamil Iskra, Ben Clifford |
SC | 2 |