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 19 records

Automated Calibration of Parallel and Distributed Computing Simulators: A Case Study

Many parallel and distributed computing research results are obtained in simulation, using simulators that mimic real-world executions on some target system. Each such simulator is configured by picking values for parameters that define the behavior of the underlying simulation models it implements. The main concern for a simulator is accuracy: simulated behaviors should be as close as possible to those observed in the real-world target system. This requires that values for each of the simulator's parameters be carefully picked, or “calibrated,” based on ground-truth real-world executions. Examining the current state of the art shows that simulator calibration, at least in the field of parallel and distributed computing, is often undocumented (and thus perhaps often not performed) and, when documented, is described as a labor-intensive, manual process. In this work we evaluate the benefit of automating simulation calibration using simple algorithms. Specifically, we use a real-world case study from the field of High Energy Physics and compare automated calibration to calibration performed by a domain scientist. Our main finding is that automated calibration is on par with or significantly outperforms the calibration performed by the domain scientist. Furthermore, automated calibration makes it straightforward to operate desirable tradeoffs between simulation accuracy and simulation speed.

Mc donald, Jesse↗

Determining Levels of Detail for Simulators of Parallel and Distributed Computing Systems via Automated Calibration

There are two sources of inaccuracy when simulating parallel and distributed computing systems: (i) a simulator implemented at an insufficient level of detail; and (ii) incorrectly calibrated simulation parameter values. Increasing the simulator’s level of detail can improve accuracy, but at the cost of higher space, time, and/or software complexity. Furthermore, evaluating the intrinsic accuracy of a simulator requires that its parameters be well-calibrated. Making decisions regarding the level of detail is thus challenging. We propose a methodology for instantiating the simulation calibration process and a framework for automating this process, which makes it possible to pick appropriate levels of detail for any simulator. We demonstrate the usefulness of our approach via two case studies for two different domains.

McDonald, Jessie [University of Hawaii at Manoa, H↗

International ACM Symposium on High Performance Parallel and Distributed Computing Conference for 2017, 2018, 2019, and 2020

The 28th ACM HPDC Conference was held in Phoenix, Arizona, June 24 and 28, 2019 (hpdc.org/2019), that was colocated with ACM FCRC 2019 (fcrc.acm.org). During the conference, Prof. Geoffrey Fox, Indiana University, was given the HPDC Achievement Award for 2019. Prof gave a keynote speech entitled “Perspectives on High-Performance Computing in a Big Data World. In addition, to the keynote speakers from HPDC and FCRC conferences, the conference organized successfully five workshops and one Ph.D. forum. The ACM FCRC had a total of 2700 attendees, and HPDC had a total of 120 attendees that included 32 students. We have used the DOE sponsorship to support the conference proceedings that acknowledge the DoE support and partially supported the travel to the HPDC PC meeting, Keynote speaker accommodation, best papers, presentation and poster award.

42 ENGINEERING↗

Contingency Analysis Based on Partitioned and Parallel Holomorphic Embedding

In the steady-state contingency analysis, the traditional Newton-Raphson method suffers from non-convergence issues when solving post-outage power flow problems, which hinders the integrity and accuracy of security assessment. In this paper, we propose a novel robust contingency analysis approach based on holomorphic embedding (HE). Here, the HE-based simulator provides theoretical convergence guarantee, which is desirable because it avoids the influence of numerical issues and provides a credible security assessment conclusion. In addition, based on the multi-area characteristics of real-world power systems, a partitioned HE (PHE) method is proposed with an interfacebased partitioning of HE formulation. The PHE method does not undermine the numerical robustness of HE and significantly reduces the computation burden in large-scale contingency analysis. The PHE method is further enhanced by parallel or distributed computation to become parallel PHE (P2HE). Tests on a 458-bus system, a synthetic 419-bus system and a large-scale 21447-bus system demonstrate the advantages of the proposed methods in robustness and efficiency.

42 ENGINEERING↗

Real‐time XFEL data analysis at SLAC and NERSC: A trial run of nascent exascale experimental data analysis

X‐ray scattering experiments using free electron lasers (XFELs) are a powerful tool to determine the molecular structure and function of unknown samples (such as COVID‐19 viral proteins). XFEL experiments are a challenge to computing in two ways: (i) due to the high cost of running XFELs, a fast turnaround time from data acquisition to data analysis is essential to make informed decisions on experimental protocols; (ii) data‐collection rates are growing exponentially, requiring new scalable algorithms. Here we report our experiences analyzing data from two experiments at the Linac Coherent Light Source (LCLS) during September 2020. Raw data were analyzed on NERSC's Cori XC40 system, using the Superfacility paradigm: our workflow automatically moves raw data between LCLS and NERSC, where it is analyzed using the software package CCTBX. We achieved real time data analysis with a turnaround time from data acquisition to full molecular reconstruction in as little as 10 min—sufficient time for the experiment's operators to make informed decisions. By hosting the data analysis on Cori, and by automating LCLS‐NERSC interoperability, we achieved a data analysis rate which matches the data acquisition rate. Completing data analysis within 10 min is a first for XFEL experiments and an important milestone if we are to keep up with data‐collection trends.

46 INSTRUMENTATION RELATED TO NUCLEAR SCIENCE AND ↗

Scalable Knowledge Graph Analytics at 136 Petaflop/s

We are motivated by newly proposed methods for data mining large-scale corpora of scholarly publications, such as the full biomedical literature, which may consist of tens of millions of papers spanning decades of research. In this setting, analysts seek to discover how concepts relate to one another. They construct graph representations from annotated text databases and then formulate the relationship-mining problem as one of computing all-pairs shortest paths (APSP), which becomes a significant bottleneck. In this context, we present a new high-performance algorithm and implementation of the Floyd-Warshall algorithm for distributed-memory parallel computers accelerated by GPUs, which we call DSNAPSHOT (Distributed Accelerated Semiring All-Pairs Shortest Path). For our largest experiments, we ran DSNAPSHOT on a connected input graph with millions of vertices using 4, 096nodes (24,576GPUs) of the Oak Ridge National Laboratory's Summit supercomputer system. We find DSNAPSHOT achieves a sustained performance of 136×1015 floating-point operations per second (136petaflop/s) at a parallel efficiency of 90% under weak scaling and, in absolute speed, 70% of the best possible performance given our computation (in the single-precision tropical semiring or “min-plus” algebra). Looking forward, we believe this novel capability will enable the mining of scholarly knowledge corpora when embedded and integrated into artificial intelligence-driven natural language processing workflows at scale.

Kannan, Ramakrishnan {ramki}↗

An exploration of online-simulation-driven portfolio scheduling in Workflow Management Systems

Workflow Management Systems used to automate the execution of scientific workflow applications on parallel and distributed computing platforms must make scheduling decisions at runtime. A large number of workflow scheduling algorithms have been proposed in the literature, but often these algorithms are evaluated based on simplifying assumptions that may not hold in practice. Furthermore, published algorithm evaluation and/or comparison results are necessarily only for a subset of all possible scenarios, and thus may not include scenarios relevant to particular use-cases. Consequently, it is difficult for Workflow Management Systems (WMSs) developers to decide which scheduling algorithm should be implemented. To obviate this difficulty, one possible approach is to implement a portfolio of scheduling algorithms and select the most effective algorithm at runtime. One method for performing this selection is to run an online simulation for each algorithm in the portfolio. The algorithm that leads to the best performance, in simulation, is selected for future use. The above simulation-driven portfolio scheduling (SDPS) approach has been proposed in a few parallel and distributed computing contexts. The main objective of this work is to evaluate the feasibility and potential merit of SDPS if implemented in WMSs. Here we perform this evaluation using simulated WMS executions, where the simulations are instantiated from real-world platform and workflow configurations. Our main finding is that SDPS is on par with or outperforms an approach in which a single algorithm is used, where this algorithm is the one that performs best on average across all our experimental scenarios. Furthermore, we find that SDPS remains an attractive proposition even in the presence of high levels of simulation error and for simulators with relatively low levels of sophistication. In many of our experimental scenarios we find that mitigating simulation error at runtime can further improve performance. Finally, we show that simulation overhead can be made sufficiently low for SDPS to be feasible in practice.

97 MATHEMATICS AND COMPUTING↗

GPU-acceleration of the ELPA2 distributed eigensolver for dense symmetric and hermitian eigenproblems

The solution of eigenproblems is often a key computational bottleneck that limits the tractable system size of numerical algorithms, among them electronic structure theory in chemistry and in condensed matter physics. Large eigenproblems can easily exceed the capacity of a single compute node, thus must be solved on distributed-memory parallel computers. We here present GPU-oriented optimizations of the ELPA two-stage tridiagonalization eigensolver (ELPA2). On top of cuBLAS-based GPU offloading, we add a CUDA kernel to speed up the back-transformation of eigenvectors, which can be the computationally most expensive part of the two-stage tridiagonalization algorithm. Furthermore, we benchmark the performance of this GPU-accelerated eigensolver on two hybrid CPU–GPU architectures, namely a compute cluster based on Intel Xeon Gold CPUs and NVIDIA Volta GPUs, and the Summit supercomputer based on IBM POWER9 CPUs and NVIDIA Volta GPUs. Consistent with previous benchmarks on CPU-only architectures, the GPU-accelerated two-stage solver exhibits a parallel performance superior to the one-stage counterpart. Finally, we demonstrate the performance of the GPU-accelerated eigensolver developed in this work for routine semi-local KS-DFT calculations comprising thousands of atoms.

97 MATHEMATICS AND COMPUTING↗

Avatar Tools

Supervised machine learning is the process of using past experience to predict the future. "Ensembles" are a machine-learning meta-method that can be applied to most machine learning algorithms. Ensembles generally greatly improve accuracy, reduce or remove most of the design issues presented by machine learning, and are admirably suited to parallel and distributed computation. The Avatar Tools codes are an implementation of ensembles specifically for decision trees. Some features that distinguish Avatar Tools from other "ensembles for decision trees" codes are: (1) Does the bookkeeping necessary for out of bag (OOB) validation. (2) Can use OOB validation to automatically determine optimal ensemble size. (3) Provides an MPI-based parallel implementation, for distributed operation. (4) Provides convenient tools for cross-validation, to assess the accuracy provided by a training set. SAND2020-3858 M Sandia National Laboratories is a multimission laboratory managed and operated by National Technology & Engineering Solutions of Sandia, LLC, a wholly owned subsidiary of Honeywell International Inc., for the U.S. Department of Energy’s National Nuclear Security Administration under contract DE-NA0003525.

Siefert, Christopher↗

Building an Integrated Ecosystem of Computational and Observational Facilities to Accelerate Scientific Discovery

Future scientific discoveries will rely on flexible ecosystems that incorporate modern scientific instruments, high performance computing resources, parallel distributed data storage, and performant networks across multiple, independent facilities. In addition to connecting physical resources, such an ecosystem presents many challenges in logistics and accessibility, especially in orchestrating computations and experiments that span across leadership computing systems and experimental instruments. Past efforts have typically been application-specific or limited to interfaces for computing resources. This paper proposes a general framework for integrating computation resources and instrument operations, addressing challenges in code development/execution, data staging and collection, software stack, control mechanisms, resource authorization and governance, and hardware integration. We also describe a demonstration use case wherein a Bayesian optimization algorithm running on an edge computing resource guides a scanning probe microscope to autonomously and intelligently characterize a material sample. This science edge ecosystem framework will provide a blueprint for federating multi-institutional, disparate resources and orchestrating scientific workflows across them to enable next-generation discoveries.

Somnath, Suhas↗

MILK : a Python scripting interface to MAUD for automation of Rietveld analysis

Modern diffraction experiments ( e.g. in situ parametric studies) present scientists with many diffraction patterns to analyze. Interactive analyses via graphical user interfaces tend to slow down obtaining quantitative results such as lattice parameters and phase fractions. Furthermore, Rietveld refinement strategies ( i.e. the parameter turn-on-off sequences) tend to be instrument specific or even specific to a given dataset, such that selection of strategies can become a bottleneck for efficient data analysis. Managing multi-histogram datasets such as from multi-bank neutron diffractometers or caked 2D synchrotron data presents additional challenges due to the large number of histogram-specific parameters. To overcome these challenges in the Rietveld software Material Analysis Using Diffraction ( MAUD ), the MAUD Interface Language Kit ( MILK ) is developed along with an updated text batch interface for MAUD . The open-source software MILK is computer-platform independent and is packaged as a Python library that interfaces with MAUD . Using MILK , model selection ( e.g. various texture or peak-broadening models), Rietveld parameter manipulation and distributed parallel batch computing can be performed through a high-level Python interface. A high-level interface enables analysis workflows to be easily programmed, shared and applied to large datasets, and external tools to be integrated with MAUD . Through modification to the MAUD batch interface, plot and data exports have been improved. The resulting hierarchical folders from Rietveld refinements with MILK are compatible with Cinema: Debye–Scherrer , a tool for visualizing and inspecting the results of multi-parameter analyses of large quantities of diffraction data. In this manuscript, the combined Python scripting and visualization capability of MILK is demonstrated with a quantitative texture and phase analysis of data collected at the HIPPO neutron diffractometer.

97 MATHEMATICS AND COMPUTING↗

Communication Lower Bounds and Optimal Algorithms for Symmetric Matrix Computations

In this article, we focus on the communication costs of three symmetric matrix computations: (i) multiplying a matrix with its transpose, known as a symmetric rank-k update (SYRK) (ii) adding the result of the multiplication of a matrix with the transpose of another matrix and the transpose of that result, known as a symmetric rank-2k update (SYR2K) (iii) performing matrix multiplication with a symmetric input matrix (SYMM). All three computations appear in the Level 3 Basic Linear Algebra Subroutines (BLAS) and have wide use in applications involving symmetric matrices. We establish communication lower bounds for these kernels using sequential and distributed-memory parallel computational models, and we show that our bounds are tight by presenting communication-optimal algorithms for each setting. Our lower bound proofs rely on applying a geometric inequality for symmetric computations and analytically solving constrained nonlinear optimization problems. As a result, the symmetric matrix and its corresponding computations are accessed and performed according to a triangular block partitioning scheme in the optimal algorithms.

Al Daas, Hussam [Rutherford Appleton Laboratory, D↗

UPC++ v1.0 Programmer’s Guide, Revision 2020.10.0

UPC++ is a C++11 library that provides Partitioned Global Address Space (PGAS) programming. It is designed for writing parallel programs that run efficiently and scale well on distributed-memory parallel computers. The PGAS model is single program, multiple-data (SPMD), with each separate constituent process having access to local memory as it would in C++. However, PGAS also provides access to a global address space, which is allocated in shared segments that are distributed over the processes. UPC++ provides numerous methods for accessing and using global memory. In UPC++, all operations that access remote memory are explicit, which encourages programmers to be aware of the cost of communication and data movement. Moreover, all remote-memory access operations are by default asynchronous, to enable programmers to write code that scales well even on hundreds of thousands of cores.

97 MATHEMATICS AND COMPUTING↗

UPC++ v1.0 Programmer’s Guide, Revision 2020.3.0

UPC++ is a C++11 library that provides Partitioned Global Address Space (PGAS) programming. It is designed for writing parallel programs that run efficiently and scale well on distributed-memory parallel computers. The PGAS model is single program, multiple-data (SPMD), with each separate constituent process having access to local memory as it would in C++. However, PGAS also provides access to a global address space, which is allocated in shared segments that are distributed over the processes. UPC++ provides numerous methods for accessing and using global memory. In UPC++, all operations that access remote memory are explicit, which encourages programmers to be aware of the cost of communication and data movement. Moreover, all remote-memory access operations are by default asynchronous, to enable programmers to write code that scales well even on hundreds of thousands of cores.

97 MATHEMATICS AND COMPUTING↗

Gradient-Based Multi-Area Distribution System State Estimation

The increasing distributed and renewable energy resources and controllable devices in distribution systems make fast distribution system state estimation (DSSE) crucial in system monitoring and control. We consider a large multi-phase distribution system and formulate DSSE as a weighted least squares (WLS) problem. We divide the large distribution system into smaller areas of subtree structure, and by jointly exploring the linearized power flow model and the network topology, we propose a gradient-based multi-area algorithm to exactly and efficiently solve the WLS problem. The proposed algorithm enables distributed and parallel computation of the state estimation problem without compromising any performance. Numerical results on a 4,521-node test feeder show that the designed algorithm features fast convergence and accurate estimation results. Comparison with traditional Gauss-Newton method shows that the proposed method has much better performance in distribution systems with a limited amount of reliable measurement. The real-time implementation of the algorithm tracks time-varying system states with high accuracy.

41 EE - Solar Energy Technologies Office (EE-4S)↗

UPC++ v1.0 Programmer’s Guide, Revision 2021.9.0

UPC++ is a C++ library that provides Partitioned Global Address Space (PGAS) programming. It is designed for writing parallel programs that run efficiently and scale well on distributed-memory parallel computers. The PGAS model is single program, multiple-data (SPMD), with each separate constituent process having access to local memory as it would in C++. PGAS additionally provides one-sided Remote Memory Access (RMA) to a global address space, which is allocated in shared segments that are distributed over the processes. UPC++ also features Remote Procedure Call (RPC) communication, making it easy to move computation to operate on data that resides on remote processes. In UPC++, all communication operations are explicit, which encourages programmers to be aware of the cost of communication and data movement. Moreover, all communication operations are asynchronous by default, to enable programmers to write code that scales well even on hundreds of thousands of cores.

96 KNOWLEDGE MANAGEMENT AND PRESERVATION↗

UPC++ v1.0 Programmer’s Guide (Rev. 2023.9.0)

UPC++ is a C++ library that supports Partitioned Global Address Space (PGAS) programming. It is designed for writing efficient, scalable parallel programs on distributed-memory parallel computers. The key communication facilities in UPC++ are one-sided Remote Memory Access (RMA) and Remote Procedure Call (RPC). The UPC++ control model is single program, multiple-data (SPMD), with each separate constituent process having access to local memory as it would in C++. The PGAS memory model additionally provides one-sided RMA communication to a global address space, which is allocated in shared segments that are distributed over the processes. UPC++ also features Remote Procedure Call (RPC) communication, making it easy to move computation to operate on data that resides on remote processes. UPC++ was designed to support exascale high-performance computing, and the library interfaces and implementation are focused on maximizing scalability. In UPC++, all communication operations are syntactically explicit, which encourages programmers to consider the costs associated with communication and data movement. Moreover, all communication operations are asynchronous by default, encouraging programmers to seek opportunities for overlapping communication latencies with other useful work. UPC++ provides expressive and composable abstractions designed for efficiently managing aggressive use of asynchrony in programs. Together, these design principles are intended to enable programmers to write applications using UPC++ that perform well even on hundreds of thousands of cores.

97 MATHEMATICS AND COMPUTING↗