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 127 records · Page 7

New method for detecting fast neutrino flavor conversions in core-collapse supernova models with two-moment neutrino transport

Fast-pairwise neutrino oscillations potentially affect many aspects of core-collapse supernovae (CCSNe): the explosion mechanism, neutrino signals, and nucleosynthesis in the ejecta. This particular mode of collective neutrino oscillations has a deep connection to the angular structure of neutrinos in momentum space; for instance, the appearance of electron-neutrino lepton number (ELN) angular crossings in momentum space is a good indicator of the occurrence of flavor conversions. However, many multidimensional CCSN simulations are carried out with approximate neutrino transport (such as two-moment methods), which limits access to the angular distributions of neutrinos, i.e., inhibits ELNcrossing searches. In this paper, we develop a new method of searching for ELN crossing in these CCSN simulations. The required data is the zeroth and first angular moments of neutrinos and the matter profile, all of which are available in CCSN models with the two-moment method. One of the novelties of our new method is the use of a ray-tracing neutrino transport to determine ELNs in the direction of the stellar center. It is designed to compensate for shortcomings of the crossing searches with only two angular moments. We assess the capability of the method by carrying out a detailed comparison with results of full Boltzmann neutrino transport in 1D and 2D CCSN models. In this work, we find that the ray-tracing neutrino transport improves the accuracy of crossing searches; indeed, the appearance/disappearance of the crossings is accurately detected even in the region of forward-peaked angular distributions. The new method is computationally cheap and has the benefit of efficient parallelization; hence, it will be useful for ELN-crossing searches in any CCSN models that employ two-moment neutrino transport.

79 ASTRONOMY AND ASTROPHYSICS↗

ARPA E GO Competition

A Nonlinear Programming SC-ACOPF Framework with Parallel Computing Capabilities

24 POWER TRANSMISSION AND DISTRIBUTION↗

Characterizing Output Bottlenecks of a Production Supercomputer: Analysis and Implications

This article studies the I/O write behaviors of the Titan supercomputer and its Lustre parallel file stores under production load. The results can inform the design, deployment, and configuration of file systems along with the design of I/O software in the application, operating system, and adaptive I/O libraries.We propose a statistical benchmarking methodology to measure write performance across I/O configurations, hardware settings, and system conditions. Moreover, we introduce two relative measures to quantify the write-performance behaviors of hardware components under production load. In addition to designing experiments and benchmarking on Titan, we verify the experimental results on one real application and one real application I/O kernel, XGC and HACC IO, respectively. These two are representative and widely used to address the typical I/O behaviors of applications.In summary, we find that Titan’s I/O system is variable across the machine at fine time scales. This variability has two major implications. First, stragglers lessen the benefit of coupled I/O parallelism (striping). Peak median output bandwidths are obtained with parallel writes to many independent files, with no striping or write sharing of files across clients (compute nodes). I/O parallelism is most effective when the application—or its I/O libraries—distributes the I/O load so that each target stores files for multiple clients and each client writes files on multiple targets in a balanced way with minimal contention. Second, our results suggest that the potential benefit of dynamic adaptation is limited. In particular, it is not fruitful to attempt to identify “good locations” in the machine or in the file system: component performance is driven by transient load conditions and past performance is not a useful predictor of future performance. For example, we do not observe diurnal load patterns that are predictable.

97 MATHEMATICS AND COMPUTING↗

Distributed non-negative matrix factorization with determination of the number of latent features

The holistic analysis and understanding of the latent (that is, not directly observable) variables and patterns buried in large datasets is crucial for data-driven science, decision making and emergency response. Such exploratory analyses require devising unsupervised learning methods for data mining and extraction of the latent features, and non-negative matrix factorization (NMF) is one of the prominent such methods. NMF is based on compute-intense non-convex constrained minimization, which, for large datasets requires fast and distributed algorithms. However, current parallel implementations of NMF fail to estimate the number of latent features. In practice, identifying these features is both difficult and significant for pattern recognition and latent feature analysis, especially for large dense matrices. Here, we introduce a distributed NMF algorithm coupled with distributed custom clustering followed by a stability analysis on dense data, which we call DnMFk, to determine the number of latent variables. The results on synthetic data and the classical Swimmer data set demonstrate the accuracy of model determination while scaling nearly linearly across multiple processors for large data. Further, we employ DnMFk to determine the number of hidden features from a terabyte matrix.

97 MATHEMATICS AND COMPUTING↗

Asynchronous and Load-Balanced Union-Find for Distributed and Parallel Scientific Data Visualization and Analysis

We present a novel distributed union-find algorithm that features asynchronous parallelism and k-d tree based load balancing for scalable visualization and analysis of scientific data. Applications of union-find include level set extraction and critical point tracking, but distributed union-find can suffer from high synchronization costs and imbalanced workloads across parallel processes. In this study, we prove that global synchronizations in existing distributed union-find can be eliminated without changing final results, allowing overlapped communications and computations for scalable processing. We also use a k-d tree decomposition to redistribute inputs, in order to improve workload balancing. We benchmark the scalability of our algorithm with up to 1,024 processes using both synthetic and application data. Here, we demonstrate the use of our algorithm in critical point tracking and super-level set extraction with high-speed imaging experiments and fusion plasma simulations, respectively.

97 MATHEMATICS AND COMPUTING↗

T-FSM: A Scalable Distributed Task-Based System for Frequent Subgraph Pattern Mining from a Big Graph

Finding frequent subgraph patterns in a big graph is an important problem with many applications such as classifying chemical compounds and building indexes to speed up graph queries. Since this problem is NP-hard, some recent parallel and distributed systems have been developed to accelerate the mining. However, they often have a huge memory cost, very long running time, suboptimal load balancing, poor scale-out capability, and possibly inaccurate results. In this article, we propose an efficient system called T-FSM for parallel mining of frequent subgraph patterns in a big graph. T-FSM supports a new anti-monotonic frequentness measure called Fraction-Score, which is more accurate than the widely used MNI measure. The execution engine of T-FSM supports both intra-machine parallelism and inter-machine parallelism. For intra-machine parallelism, T-FSM adopts a novel task-based execution model to ensure high multithreading concurrency, bounded memory consumption, and effective load balancing. For inter-machine parallelism, T-FSM ensures good scale-out performance with a lightweight pattern rebalancing approach that reduces workload skewness of pattern evaluations among machines. To avoid recomputing the contexts for migrated patterns, we design a novel context cache table to support concurrent and asynchronous requesting and caching of remote context data, which can timely evict and garbage collect used pattern contexts that are no longer needed to keep memory consumption bounded. Extensive experiments show that T-FSM is orders of magnitude faster than existing state-of-the-art parallel systems (more than 10×, 51×, 131×, 55× speedup over ScaleMine, DistGraph, Pangolin and Peregrine, respectively) and distributed systems (more than 42× and 88× over ScaleMine and DistGraph, respectively) for frequent subgraph pattern mining, and it scales out satisfactorily to 512 CPU cores on the Polaris supercomputer at Argonne National Laboratory.

97 MATHEMATICS AND COMPUTING↗

GridOPTICS/GridPACK

GridPACK is a software framework consisting of a set of modules designed to simplify the development of programs that model the power grid and run on parallel, high performance computing platforms. It also contains several fully developed applications, including powerflow, dynamic simulation, state estimation, Kalman filter analysis (dynamic state estimation), contingency analysis and real time path rating. These applications can be used either standalone or as components in more complicated workflows that combine several different types of application together. The framework modules are available as a combination of libraries and software templates and consist of components for setting up and distributing power grid networks, support for modeling the behavior of individual buses and branches in the network, converting the network models to the corresponding algebraic equations, and parallel routines for manipulating and solving large algebraic systems. The framework also contains a module for distributing tasks evenly amongst computing resources, even if individual tasks vary widely in their execution times. Additional modules support input and output, basic statistical analysis of contingency based calculations, distributed data structures, as well as basic profiling and error management.

Palmer, Bruce↗

Scalable Computation of Topological Abstractions for Scalar Data

Topological data analysis has become an important tool for large scale scalar data analysis and visualization, efficiently extracting the inherent structure and features of interest of the data. However, with growing dataset sizes and complexity, it is increasingly becoming infeasible to compute topological abstractions of interest in serial and on single machines. This paper presents the state of the art in the scalable computation of topological abstractions on scalar data, in shared memory parallel on single machines, and in distributed memory parallel on multiple machines. We highlight results for set‐based, graph‐based and complex‐based abstractions and organize the state of the art based on this taxonomy. The paper identifies parallelization and distribution techniques common in topological algorithms and highlights further areas of interest with underdeveloped efforts.

97 MATHEMATICS AND COMPUTING↗

Distributed Macroscopic Traffic Simulation with Open Traffic Models

This paper presents OTM-MPI, an extension of the Open Traffic Models platform (OTM) for running macroscopic traffic simulations in high-performance computing environments. OTM-MPI represents the first open-source, distributed-memory, macroscopic simulation model developed for modern high performance parallel machines and large networks. Macroscopic simulations are appropriate for studying regional traffic scenarios when aggregate trends are of interest, rather than individual vehicle traces. They are also appropriate for studying the routing behavior of classes of vehicles, such as app-informed vehicles. The network partitioning was performed with METIS. Inter-process communication was done with MPI (message-passing interface). Results are provided for two networks: one realistic network which was obtained from Open Street Maps for Chattanooga, TN, and another larger synthetic grid network. The software recorded a speedups of 198x using 256 cores for Chattanooga, and 475x with 1,024 cores for the synthetic network.

macro-scopic traffic simulation↗

TriC: Distributed-memory Triangle Counting by Exploiting the Graph Structure

Graph analytics has emerged as an important tool in the analysis of large scale data from diverse application domains such as social networks, cyber security and bioinformatics. Counting the number of triangles in a graph is a fundamental kernel with several applications such as detecting the community structure of a graph or in identifying important vertices in a graph. The ubiquity of massive datasets is driving the need to scale graph analytics on parallel systems. However, numerous challenges exist in efficiently parallelizing graph algorithms, especially on distributed-memory systems. Irregular memory accesses and communication patterns, low computation to communication ratios, and the need for frequent synchronization are some of the leading challenges. In this paper, we present TriC, our distributed-memory implementation of triangle counting in graphs using the Message Passing Interface (MPI), as a submission to the 2020 GraphChallenge competition. Using a set of synthetic and real-world inputs from the challenge, we demonstrate a speedup of up to 90x relative to previous work on 32 processor-cores of a NERSC Cori node. We also provide details from distributed runs with up to8192 processes along with strong scaling results. The observations presented in this work provide an understanding of the system-level bottlenecks at scale that specifically impact sparse-irregular workloads and will therefore benefit other efforts to parallelize graph algorithms.

Halappanavar, Mahantesh↗

QCOR; A Language Extension Specification for the Heterogeneous Quantum-Classical Model of Computation

Quantum computing (QC) is an emerging computational paradigm that leverages the laws of quantum mechanics to perform elementary logic operations. Existing programming models for QC were designed with fault-tolerant hardware in mind, envisioning stand-alone applications. However, the susceptibility of near-term quantum computers to noise limits their stand-alone utility. To better leverage limited computational strengths of noisy quantum devices, hybrid algorithms have been suggested whereby quantum computers are used in tandem with their classical counterparts in a heterogeneous fashion. This modus operandi calls out for a programming model and a high-level programming language that natively and seamlessly supports heterogeneous quantum-classical hardware architectures in a single-source-code paradigm. Motivated by the lack of such a model, we introduce a language extension specification, called QCOR, which enables single-source quantum-classical programming. Programs written using the QCOR library–based language extensions can be compiled to produce functional hybrid binary executables. After defining QCOR’s programming model, memory model, and execution model, we discuss how QCOR enables variational, iterative, and feed-forward QC. Additionally, QCOR approaches quantum-classical computation in a hardware-agnostic heterogeneous fashion and strives to build on best practices of high-performance computing. The high level of abstraction in the language extension is intended to accelerate the adoption of QC by researchers familiar with classical high-performance computing.

97 MATHEMATICS AND COMPUTING↗

An Open-Source Parallel EMT Simulation Framework: Preprint

As the integration level of inverter-based resources (IBR) increases, ensuring the reliable operation of the bulk power systems requires the use of electromagnetic transient (EMT) simulation tools to identify and mitigate system-wide stability risks. Conducting EMT studies for large-scale, IBR-rich grids, however, is challenging due to the inherent computational bottleneck caused by the underlying high-fidelity models and required small time steps. This paper introduces ParaEMT: an open-source, generic EMT simulation framework designed to accelerate simulations by leveraging advanced parallel computational technologies, such as high-performance computers. This paper presents a comprehensive exposition of ParaEMT, covering its modeling library, simulation strategy, framework structure, operational procedures, and auxiliary features, alongside its extensible parallel computational architecture. Notably, ParaEMT is a publicly accessible and modularized framework written in Python, thereby facilitating future development and the integration of new models and algorithms. The accuracy and efficiency of ParaEMT are demonstrated by rigorous validations via multiple case studies.

electromagnetic transient simulation↗

Toucan: A performance portable, scalable implementation of the DECA algorithm

In the field of additive manufacturing (AM), cellular automata (CA) is extensively used to simulate microstructural evolution during solidification. However, while traditional CA approaches are relatively fast, they still require a substantial number of time steps, are limited to moderate volumes, and are relatively difficult to improve through parallelism due to the highly localized nature of the solidification front. Here, to address these issues of time to solution and load balancing, we introduce Toucan, a parallel, performance-portable, and scalable code written in C++ with the Kokkos library that leverages the discrete event inspired cellular automata (DECA) algorithm to perform parallel-in-time (PinT) grain growth simulations. Toucan effectively mitigates load balancing issues by distributing the computational workload more evenly across processors, enhancing scalability and efficiency. We conduct both strong and weak scaling studies on up to 64 GPUs on the Frontier supercomputer, demonstrating that Toucan significantly outperforms the current state-of-the-art, time-stepped CA code, ExaCA, on both single and multi-GPU simulations. Even in AM-specific weak scaling scenarios, Toucan maintains near-ideal scaling, in contrast to the linear increase observed with ExaCA due to the moving laser raster pattern. This study highlights Toucan’s potential to transform microstructural simulations in AM by radically improving both efficiency and scalability over existing methods.

36 MATERIALS SCIENCE↗

Distributed Multi-GPU Community Detection on Exascale Computing Platforms

Community detection is a fundamental operation in graph mining, and by uncovering hidden structures and patterns within complex systems it helps solve fundamental problems pertaining to social networks, such as information diffusion, epidemics, and recommender systems. Scaling graph algorithms for massive networks becomes challenging on modern distributed-memory multi-GPU (Graphics Processing Unit) systems due to limitations such as irregular memory access patterns, load imbalances, higher communication-computation ratios, and cross-platform support. We present a novel algorithm HiPDPL-GPU (distributed parallel Louvain) to address these challenges. We conduct experiments involving different partitioning techniques to achieve optimized performance of HiPDPL-GPU on the two largest supercomputers: Frontier and Summit. Remarkably, HiPDPL-GPU processes a graph with 4.2 billion edges in less than 3 minutes using 1024 GPUs. Qualitatively performance of HiPDPL-GPU is similar or better compared to other state-of-the-art CPU- and GPU-based implementations. While prior GPU implementations have predominantly employed CUDA, our first-of-its-kind implementation for community detection is cross-platform, accommodating both AMD and NVIDIA GPUs.

graph algorithms, high performance comptuing↗

Distributed Multi-GPU Community Detection on Exascale Computing Platforms

Community detection is a fundamental operation in graph mining, and by uncovering hidden structures and patterns within complex systems it helps solve fundamental problems pertaining to social networks, such as information diffusion, epidemics, and recommender systems. Scaling graph algorithms for massive networks becomes challenging on modern distributed-memory multi-GPU (Graphics Processing Unit) systems due to limitations such as irregular memory access patterns, load imbalances, higher communication-computation ratios, and cross-platform support. We present a novel algorithm HiPDPL-GPU (Distributed Parallel Louvain) to address these challenges. We conduct experiments involving different partitioning techniques to achieve an optimized performance of HiPDPL-GPU on the two largest supercomputers: Frontier and Summit. Remarkably, HiPDPL-GPU processes a graph with 4.2 billion edges in less than 3 minutes using 1024 GPUs. Qualitatively, the performance of HiPDPL-GPU is similar or better compared to other state-of-the-art CPU- and GPU-based implementations. While prior GPU implementations have predominantly employed CUDA, our first-of-its-kind implementation for community detection is cross-platform, accommodating both AMD and NVIDIA GPUs.

Sattar, Naw Safrin↗

libEnsemble: A complete Python toolkit for dynamic ensembles of calculations

Almost all science and engineering applications eventually stop scaling: their runtime no longer decreases as available computational resources increase. Therefore, many applications will struggle to efficiently use emerging extreme-scale high-performance, parallel, and distributed systems. libEnsemble is a complete Python toolkit and workflow system for intelligently driving ensembles of experiments or simulations at massive scales. It enables and encourages multidisciplinary design, decision, and inference studies portably running on laptops, clusters, and supercomputers.

97 MATHEMATICS AND COMPUTING↗

Situational Awareness of Grid Anomalies (SAGA)

The modern power industry becomes more vulnerable to cyber events due to the growing interconnectivity, interdependence, and complexity of the electric power grid. High-fidelity modeling and simulation tools that support the preventative risk analysis on potential cyber-relevant events is essential for ensuring the situational awareness of the system operator as it provides an inexpensive and risk-free environment to test the system responses under various cyber-relevant events and hereby can support research on cyber anomaly detection, optimal protective resource allocation, and mitigation measures. In this webinar, we will share NREL's cybersecurity research capabilities by highlighting the development of a scalable cyber-physical event test bed and demonstration with real hardware in the loop. The developed cyber-physical event test bed is backboned by an integrated transmission, distribution, and communication dynamic co-simulation framework and a plug-and-play cyber event generation module. It is designed to be modular and compatible with parallel computing, and thereby supports large-scale system simulations at an affordable computation cost. The test bed can capture millisecond-to-minutes dynamic frequency and voltage responses under cyber events from the bulk transmission system to the active distribution systems and distributed energy resources at the grid edge.

co-simulation↗