Engineering Papers⌕ Search

SEARCH · Engineering Papers

Results for “parallel and distributed computing”

Search indexed NASA NTRS and DOE OSTI research on propulsion, heat transfer, battery materials and energy systems. Follow report and document links to the original sources.

Quote a phrase for an exact phrase match. Source license links do not imply unrestricted reuse.

At least 181 records · Page 10

Machine Learning for Distributed Acoustic Sensing data (MLDAS) v1.0.1

MLDAS is a Python-written package for exploratory data analysis and deep learning training on Distributed Acoustic Sensing data. The machine learning tools are powered by the PyTorch library and designed to work efficiently on large scale datasets using parallel computing. Various SLURM scripts as well as a tutorial have also been made available to allow geophysicists to quickly and easily implement the available tools in their analysis workflow on supercomputer facilities.

Dumont, Vincent↗

SBIR Phase I Final Report, TACO: Distributed and Heterogeneous Sparse Compiler

Tensor algebra is a powerful tool for computing, but writing optimized codes that operate on sparse tensors can be very complex. This project enables a Tensor Algebra Compiler (TACO) that simplifies this task from man-years to man-days and extends TACO to support complex and large distributed systems. This report details the hypotheses, approaches used, and findings in this project.

97 MATHEMATICS AND COMPUTING↗

Pele: An Exascale-Ready Suite of Combustion Codes

High fidelity simulations of realistic combustion devices are extremely demanding computationally because of the requirements to capture complex fuel chemical decomposition, its intricate interactions with turbulent, often multiphase, flows, and the wide separation of space and time scales between the thin flame and the device boundaries. Software required to carry out such computations tends to be extremely complex, particularly when designed to exploit hardware accelerators, and can be difficult to port and maintain. We present Pele, a performance portable suite of tools for the simulation of combustion systems, including codes to evolve reactive multiphase configurations in the low Mach number and compressible flow regimes, along with a set of inter-compatible post processing and in situ analysis tools. The Pele suite of tools is built on top of the AMReX framework for block-structured adaptive mesh refinement, which provides efficient data structures and algorithms that enable the development of a wide variety of efficient mesh and particle based PDE integration schemes. A hierarchical MPI+X parallelism scheme supports CPU-only and accelerated architectures, where X can be OpenMP, CUDA, and HIP based approaches for intra-node computational work distribution. The algorithms and data structures underlying the Pele simulation and analysis tools are highly scalable and performant across a wide variety of high-performance computing platforms, including DOEs newest exascale-class machines, Frontier and Aurora. The simulation and analysis tools are fully documented and freely distributed as open source via GitHub. We present key algorithmic and software challenges, solution strategies, performance and resulting set of capabilities.

AMReX↗

Scalable quantum computational science: A perspective from block-encodings and polynomial transformations

Significant developments made in quantum hardware and error correction recently have been driving quantum computing toward practical utility. However, gaps remain between abstract quantum algorithmic development and practical applications in computational sciences. In this perspective article, we propose several properties that scalable quantum computational science methods should possess. We further discuss how block-encodings and polynomial transformations can potentially serve as a unified framework with the desired properties. Recent advancements on these topics are presented, including the construction and assembly of block-encodings, and various generalizations of quantum signal processing (QSP) algorithms to perform polynomial transformations. The scalability of QSP methods on parallel and distributed quantum architectures is also highlighted. Promising applications in simulation and observable estimation in chemistry, physics, and optimization problems are presented. We hope this perspective serves as a gentle introduction to state-of-the-art quantum algorithms for the computational science community and inspires future development of scalable quantum computational science methodologies that bridge theory and practice.

Bayesian inference↗

FitCache: A Transparent Drop-In Framework for Multi-Tier Caching to Accelerate Distributed Deep Learning Workloads

Training in Deep learning (DL) remains highly compute- and data-intensive, with I/O becoming a critical bottleneck as models and datasets scale. Recent studies report that data loading can dominate training time, especially on large-scale HPC systems with shared parallel file systems (PFS). Existing caching approaches either rely on single-tier designs or require intrusive modifications to training pipelines, limiting their portability and effectiveness. In this work, we present FitCache, a transparent drop-in framework for multi-tier caching to accelerate distributed DL training by coordinating fast local memory (e.g., DRAM, Persistent Memory (PMem)) and NVMe as hierarchical caches atop PFS. Our design adapts to hardware diversity, i.e., if NVMe is missing, memory transparently acts as a caching tier, ensuring stable performance. FitCache transparently intercepts I/O requests and issues concurrent fetches across all tiers, returning data from the fastest responder without centralized metadata or static redirection paths. FitCache adapts to dynamic workloads and heterogeneous clusters while maintaining POSIX compatibility. Experiments on Frontier (2048 GPUs) and smaller research clusters show that FitCache reduces training time by up to 40% and per-batch I/O latency by up to 71.6% compared to Lustre Orion PFS, offering a drop-in solution for scalable DL training.

Hu, Guangxing [ORNL] (ORCID:0009000283203614)↗

Episodic Earthquake Swarms in the Mineral Mountains, Utah Driven by the Roosevelt Hydrothermal System

Over 1,000 earthquakes (-2.0 < M < 2.0), identified using a matched-filter method, occurred in the Mineral Mountains, Utah between 2016 and 2019. The enhanced catalog is complete down to M -0.9 and contains roughly 15 times more events than originally cataloged. Earthquake relocation of ~800 earthquakes shows that activity is concentrated in a <2 km long E-W striking narrow zone, ~4 km east of the Roosevelt hydrothermal system. Two fault orientations, both N-S and E-W parallel to the Opal Mound and Mag Lee faults, respectively, are observed after computing composite focal mechanisms of highly similar earthquakes. Looking solely at the temporal distribution of the seismicity, we identify 15 periods of swarm-like activity, with two major clusters occurring in December 2016, recorded by three stations, and in October 2019 recorded by eight stations. The October 2019 swarm, the best recorded sequence in the area, provides evidence for the underlying triggering mechanism. We show that a complex mechanism of fluid diffusion and aseismic slip is responsible for the swarm evolution with migration velocities reaching 10 km/day. We hypothesize that these episodic swarms in the Mineral Mountains are primarily driven by migrating fluids that originate within the Roosevelt hydrothermal system.

58 GEOSCIENCES↗

Scalable training of graph convolutional neural networks for fast and accurate predictions of HOMO-LUMO gap in molecules

Abstract Graph Convolutional Neural Network (GCNN) is a popular class of deep learning (DL) models in material science to predict material properties from the graph representation of molecular structures. Training an accurate and comprehensive GCNN surrogate for molecular design requires large-scale graph datasets and is usually a time-consuming process. Recent advances in GPUs and distributed computing open a path to reduce the computational cost for GCNN training effectively. However, efficient utilization of high performance computing (HPC) resources for training requires simultaneously optimizing large-scale data management and scalable stochastic batched optimization techniques. In this work, we focus on building GCNN models on HPC systems to predict material properties of millions of molecules. We use HydraGNN, our in-house library for large-scale GCNN training, leveraging distributed data parallelism in PyTorch. We use ADIOS, a high-performance data management framework for efficient storage and reading of large molecular graph data. We perform parallel training on two open-source large-scale graph datasets to build a GCNN predictor for an important quantum property known as the HOMO-LUMO gap. We measure the scalability, accuracy, and convergence of our approach on two DOE supercomputers: the Summit supercomputer at the Oak Ridge Leadership Computing Facility (OLCF) and the Perlmutter system at the National Energy Research Scientific Computing Center (NERSC). We present our experimental results with HydraGNN showing (i) reduction of data loading time up to 4.2 times compared with a conventional method and (ii) linear scaling performance for training up to 1024 GPUs on both Summit and Perlmutter.

37 INORGANIC, ORGANIC, PHYSICAL, AND ANALYTICAL CH↗

A Massively Parallel Implementation of the CCSD(T) Method Using the Resolution-of-the-Identity Approximation and a Hybrid Distributed/Shared Memory Parallelization Model

In this work, a parallel algorithm is described for the coupled-cluster singles and doubles method augmented with a perturbative correction for triple excitations [CCSD(T)] using the resolution-of-the-identity (RI) approximation for two-electron repulsion integrals (ERIs). The algorithm bypasses the storage of four-center ERIs by adopting an integral-direct strategy. The CCSD amplitude equations are given in a compact quasi-linear form by factorizing them in terms of amplitude-dressed three-center intermediates. A hybrid MPI/OpenMP parallelization scheme is employed, which uses the OpenMP-based shared memory model for intranode parallelization and the MPI-based distributed memory model for internode parallelization. Parallel efficiency has been optimized for all terms in the CCSD amplitude equations. Two different algorithms have been implemented for the rate-limiting terms in the CCSD amplitude equations that entail and -scaling computational costs, where N O and N V denote the number of correlated occupied and virtual orbitals, respectively. One of the algorithms assembles the four-center ERIs requiring N V 4 and N O 2 N V 2 -scaling memory costs in a distributed manner on a number of MPI ranks, while the other algorithm completely bypasses the assembling of quartic memory-scaling ERIs and thus largely reduces the memory demand. It is demonstrated that the former memory-expensive algorithm is faster on a few hundred cores, while the latter memory-economic algorithm shows a better strong scaling in the limit of a few thousand cores. The program is shown to exhibit a near-linear scaling, in particular for the compute-intensive triples correction step, on up to 8000 cores. The performance of the program is demonstrated via calculations involving molecules with 24–51 atoms and up to 1624 atomic basis functions. As the first application, the complete basis set (CBS) limit for the interaction energy of the π-stacked uracil dimer from the S66 data set has been investigated. This work reports the first calculation of the interaction energy at the CCSD(T)/aug-cc-pVQZ level without local orbital approximation. The CBS limit for the CCSD correlation contribution to the interaction energy was found to be -8.01 kcal/mol, which agrees very well with the value -7.99 kcal/mol reported by Schmitz, Hättig, and Tew [ Phys. Chem. Chem. Phys. 2014 , 16 , 22167-22178]. The CBS limit for the total interaction energy was estimated to be -9.64 kcal/mol.

37 INORGANIC, ORGANIC, PHYSICAL, AND ANALYTICAL CH↗

Distributed out-of-memory NMF on CPU/GPU architectures

We propose an efficient distributed out-of-memory implementation of the non-negative matrix factorization (NMF) algorithm for heterogeneous high-performance-computing systems. The proposed implementation is based on prior work on NMFk, which can perform automatic model selection and extract latent variables and patterns from data. In this work, we extend NMFk by adding support for dense and sparse matrix operation on multi-node, multi-GPU systems. The resulting algorithm is optimized for out-of-memory problems where the memory required to factorize a given matrix is greater than the available GPU memory. Memory complexity is reduced by batching/tiling strategies, and sparse and dense matrix operations are significantly accelerated with GPU cores (or tensor cores when available). Input/output latency associated with batch copies between host and device is hidden using CUDA streams to overlap data transfers and compute asynchronously, and latency associated with collective communications (both intra-node and inter-node) is reduced using optimized NVIDIA Collective Communication Library (NCCL) based communicators. Benchmark results show significant improvement, from 32X to 76x speedup, with the new implementation using GPUs over the CPU-based NMFk. Good weak scaling was demonstrated on up to 4096 multi-GPU cluster nodes with approximately 25,000 GPUs when decomposing a dense 340 Terabyte-size matrix and an 11 Exabyte-size sparse matrix of density 10 -6 .

97 MATHEMATICS AND COMPUTING↗

A fast particle-based approach for calibrating a 3-D model of the Antarctic ice sheet

We consider the scientifically challenging and policy-relevant task of understanding the past and projecting the future dynamics of the Antarctic ice sheet. The Antarctic ice sheet has shown a highly nonlinear threshold response to past climate forcings. Triggering such a threshold response through anthropogenic greenhouse gas emissions would drive drastic and potentially fast sea level rise with important implications for coastal flood risks. Previous studies have combined information from ice sheet models and observations to calibrate model parameters. These studies have broken important new ground but have either adopted simple ice sheet models or have limited the number of parameters to allow for the use of more complex models. These limitations are largely due to the computational challenges posed by calibration as models become more computationally intensive or when the number of parameters increases. Here, we propose a method to alleviate this problem: a fast sequential Monte Carlo method that takes advantage of the massive parallelization afforded by modern high-performance computing systems. We use simulated examples to demonstrate how our sample-based approach provides accurate approximations to the posterior distributions of the calibrated parameters. The drastic reduction in computational times enables us to provide new insights into important scientific questions, for example, the impact of Pliocene era data and prior parameter information on sea level projections. These studies would be computationally prohibitive with other computational approaches for calibration such as Markov chain Monte Carlo or emulation-based methods. We also find considerable differences in the distributions of sea level projections when we account for a larger number of uncertain parameters. For example, based on the same ice sheet model and data set, the 99th percentile of the Antarctic ice sheet contribution to sea level rise in 2300 increases from 6.5 m to 13.1 m when we increase the number of calibrated parameters from three to 11. With previous calibration methods, it would be challenging to go beyond five parameters. Here, this work provides an important next step toward improving the uncertainty quantification of complex, computationally intensive and decision-relevant models.

54 ENVIRONMENTAL SCIENCES↗

Traveler: Navigating Task Parallel Traces for Performance Analysis

Understanding the behavior of software in execution is a key step in identifying and fixing performance issues. This is especially important in high performance computing contexts where even minor performance tweaks can translate into large savings in terms of computational resource use. To aid performance analysis, developers may collect an execution trace —a chronological log of program activity during execution. As traces represent the full history, developers can discover a wide array of possibly previously unknown performance issues, making them an important artifact for exploratory performance analysis. However, interactive trace visualization is difficult due to issues of data size and complexity of meaning. Traces represent nanosecond-level events across many parallel processes, meaning the collected data is often large and difficult to explore. The rise of asynchronous task parallel programming paradigms complicates the relation between events and their probable cause. Here, to address these challenges, we conduct a continuing design study in collaboration with high performance computing researchers. We develop diverse and hierarchical ways to navigate and represent execution trace data in support of their trace analysis tasks. Through an iterative design process, we developed Traveler , an integrated visualization platform for task parallel traces. Traveler provides multiple linked interfaces to help navigate trace data from multiple contexts. We evaluate the utility of Traveler through feedback from users and a case study, finding that integrating multiple modes of navigation in our design supported performance analysis tasks and led to the discovery of previously unknown behavior in a distributed array library.

97 MATHEMATICS AND COMPUTING↗

Asynchronous Iterative Solvers for Extreme-Scale Computing

The Asynchronous Iterative Solvers for Extreme-Scale Computing (AsyncIS) project aims to explore more efficient numerical algorithms by decreasing their overhead. AsyncIS does this by replacing the outer Krylov subspace solver with an asynchronous optimized Schwarz method, thereby removing the global synchronization and bulk synchronous operations typically used in numerical codes. AsyncIS—a U.S. Department of Energy (DOE)-funded collaboration between Georgia Tech, the University of Tennessee, Knoxville, Temple University, and Sandia National Laboratories—also focuses on the development and optimization of asynchronous preconditioners (i.e., preconditioners that are generated and/or applied in an asynchronous fashion). The novel preconditioning algorithms that provide fine-grained parallelism enable preconditioned Krylov solvers to run efficiently on large-scale distributed systems and manycore accelerators like GPUs.

97 MATHEMATICS AND COMPUTING↗

cuTS: Scaling Subgraph Isomorphism on Distributed Multi-GPUSystems Using Trie Based Data Structure

Subgraph isomorphism is a pattern-matching algorithm widely used in many domains such as chem-informatics, bioinformatics, databases, and social network analysis. It is computationally expensive and is a proven NP-hard problem. The massive parallelism offered by the GPU hardware is well suited for solving the subgraph isomorphism. However, current GPU implementations are far from the achievable performance. Moreover, the enormous memory requirement of current approaches limits the problem size that can be handled. This work analyzes the fundamental challenges associated with processing the subgraph isomorphism on GPUs and develops an efficient GPU hardware-aware implementation. We also develop a new GPU-friendly trie-based data structure to drastically reduce the intermediate storage space requirement. Hence, our approach runs larger benchmarks than the competitors. We also develop the first distributed sub-graph isomorphism algorithm for GPUs. Our experimental evaluation section demonstrates the efficacy of our approach by comparing the execution time and number of cases that we can handle against the state-of-the-art GPU implementations.

Xiang, Lizhi↗

Optimizing High Performance Markov Clustering for Pre-Exascale Architectures

HipMCL is a high-performance distributed memory implementation of the popular Markov Cluster Algorithm (MCL) and can cluster large-scale networks within hours using a few thousand CPU-equipped nodes. It relies on sparse matrix computations and heavily makes use of the sparse matrix-sparse matrix multiplication kernel (SpGEMM). The existing parallel algorithms in HipMCL are not scalable to Exascale architectures, both due to their communication costs dominating the runtime at large concurrencies and also due to their inability to take advantage of accelerators that are increasingly popular. In this work, we systematically remove scalability and performance bottlenecks of HipMCL. We enable GPUs by performing the expensive expansion phase of the MCL algorithm on GPU. Additionally, we propose a CPU-GPU joint distributed SpGEMM algorithm called pipelined Sparse SUMMA and integrate a probabilistic memory requirement estimator that is fast and accurate. Furthermore, we develop a new merging algorithm for the incremental processing of partial results produced by the GPUs, which improves the overlap efficiency and the peak memory usage. We also integrate a recent and faster algorithm for performing SpGEMM on CPUs. We validate our new algorithms and optimizations with extensive evaluations. With the enabling of the GPUs and integration of new algorithms, HipMCL is up to 12.4x faster, being able to cluster a network with 70 million proteins and 68 billion connections just under 15 minutes using 1024 nodes of ORNL's Summit supercomputer.

97 MATHEMATICS AND COMPUTING↗

Scaling the SciDAC QuantOm Workflow

As part of the Scientific Discovery through Advanced Computing (SciDAC) program, the Quantum Chromodynamics Nuclear Tomography (QuantOM) project aims to analyze data from Deep Inelastic Scattering (DIS) experiments conducted at Jefferson Lab and the upcoming Electron Ion Collider. The DIS data analysis is performed on an event-level by combining the input from theoretical and experimental nuclear physics into a single, composable workflow. The optimization itself (I.e. fitting the experimental data with theoretical predictions) is carried out by a machine / deep learning algorithm. The size of the acquired DIS data as well as the complexity of the workflow itself require that the analysis is performed across multiple GPUs on high performance computing systems, such as Polaris at Argonne National Laboratory. This presentation discusses the novelties and challenges that came along with parallelizing this workflow. Recent results are compared to common distributed training techniques.

Lersch, Daniel↗

Reinforcement Learning for Load-balanced Parallel Particle Tracing

We explore an online reinforcement learning (RL) paradigm to dynamically optimize parallel particle tracing performance in distributed-memory systems. Our method combines three novel components: (1) a work donation algorithm, (2) a high-order workload estimation model, and (3) a communication cost model. First, we design an RL-based work donation algorithm. Our algorithm monitors workloads of processes and creates RL agents to donate data blocks and particles from high-workload processes to low-workload processes to minimize program execution time. The agents learn the donation strategy on the fly based on reward and cost functions designed to consider processes' workload changes and data transfer costs of donation actions. Second, we propose a workload estimation model, helping RL agents estimate the workload distribution of processes in future computations. Third, we design a communication cost model that considers both block and particle data exchange costs, helping RL agents make effective decisions with minimized communication costs. We demonstrate that our algorithm adapts to different flow behaviors in large-scale fluid dynamics, ocean, and weather simulation data. Our algorithm improves parallel particle tracing performance in terms of parallel efficiency, load balance, and costs of I/O and communication for evaluations with up to 16,384 processors.

Distributed and parallel particle tracing↗

The Kokkos Ecosystem [Brief]

In 2016/2017, the field of High-Performance Computing (HPC) entered a new era driven by fundamental physics challenges to produce ever more energy and cost-efficient processors. Since the convergence on the Message-Passing Interface (MPI) standard in the mid-1990s, application developers enjoyed a seemingly static view of the underlying machine — that of a distributed collection of homogeneous nodes executing in collaboration. However, after almost two decades of dominance, the sole use of MPI to derive parallelism acted as a limiter to improved future performance. While MPI is widely expected to continue to function as the basic mechanism for communication between compute nodes for the immediate future, additional parallelism is required on the computing node itself if high performance and efficiency goals are to be realized. When reviewing the architectures of the top HPC systems today, the change in paradigm is clear: the compute nodes of the leading machines in the world are either powered by many-core chips with a few dozen cores each, or use heterogeneous designs, where traditional CPUs marshal work to massively parallel compute accelerators which has as many as 200,000 processing threads in flight simultaneously. Complicating matters further for application developers, each processor vendor has its own preferred way of writing code for their architecture.The Kokkos EcoSystem was released by Sandia in 2017 to address this new era in HPC system design by providing a vendor independent performance portable programming system for scientific, engineering, and mathematical software applications written in the C++ programming language. Using Kokkos, application developers can be more productive because they will not have to create and maintain separate versions of their software for each architecture, nor will they have to be experts in each architecture's peculiar requirements. Instead, they will have a single method of programming for the diverse set of modern HPC architectures. While Kokkos started in 2011 as a programming model only, it soon became clear that complex applications needed more. It is also critical to have a portable mathematical functions and developers need tools to debug their applications, gain insight into the performance characteristics of their codes and tune algorithm performance parameters through automated processes. The Kokkos EcoSystem addresses those needs through its three main components: the Kokkos Core programming model, the Kokkos Kernels math library, and the Kokkos Tools project.

97 MATHEMATICS AND COMPUTING↗

A GPU ‐Accelerated 3D Unstructured Mesh Based Particle Tracking Code for Multi‐Species Impurity Transport Simulation in Fusion Tokamaks

ABSTRACT This paper presents the multi‐species global impurity transport capability developed in a GPU‐accelerated fully 3D unstructured mesh‐based code, GITRm, to simultaneously track multiple impurity species and handle interactions of these impurities with mixed‐material surfaces. Different computational approaches to model particle‐surface interaction or surface response have been developed and compared. Sheath electric field is taken into account by employing a fast distance‐to‐boundary calculation, which is carried out in parallel on distributed or partitioned meshes on multiple GPUs without the need for any inter‐process communication during the simulation. Several example cases, including two for the DIII‐D tokamak, that is, one with the SAS‐V divertor and the other with the collector probes, are used to demonstrate the utility of the current multi‐species capability. For the DIII‐D probe case, the capability of GITRm to resolve the spatial distribution of particles in localized regions, such as diagnostic probes, within non‐axisymmetric tokamak geometries is demonstrated. These simulations involve up to 320 million particles and utilize up to 48 GPUs.

Nath, Dhyanjyoti D. [Scientific Computation Resear↗