Modern supercomputers feature an ever-increasing degree of parallelism, particularly in the number of cores per node. These high core counts are considered in our flexible implementation of allreduce, which was implemented specifically with shared-memory communication in mind. At a high level, our algorithm consists of a reduce_scatter stage followed by an allgather stage, or a reduce stage followed by a broadcast stage, and allows for different factors (aka multi-radix) to be applied at each. The reduce and broadcast operations are also considered as standalone functions. Where barriers are required, they are integrated into the algorithm using counters to track progress. To accommodate the complexity of this approach, our implementation is split into a setup phase and an execution phase. The setup phase occurs only once for a given set of parameters, and is responsible for determining the algorithm that will be run each time the allreduce is called within the execution phase. We present two interfaces: a persistent collective interface (an MPI 4.0 feature), which inherently aligns with our setup and execution phases, and a blocking interface, where the setup is performed the first time a specific algorithm is required. Using these methods, we achieve speedups of half an order of magnitude compared to the persistent and blocking allreduce implementations of MPICH and Open MPI on a dual-socket node with AMD EPYC processors, and almost an order of magnitude on a four-socket node with NVIDIA Grace-Hopper processors. Reductions on vectors residing on both CPU and GPU memory are performed. Our implementation also achieves good performance on multiple nodes. A standard benchmark of the application CP2K is sped up by 2.5%. Notably, for long messages, our implementation achieves the same performance as NCCL.
Modern supercomputers feature an ever-increasing degree of parallelism, especially in the number of cores per node. These high core counts are considered in our flexible implementation of persistent allreduce (an MPI-4 feature), which was implemented specifically with shared-memory communication in mind. At a high level, our algorithm consists of a reduce scatter stage followed by an allgather stage, and allows for different factors (i.e., multi-radix) to be applied at each. Where barriers are required, they are integrated into the algorithm using counters to track progress. In order to accommodate the complexity of this approach, our implementation is split into a setup phase and an execution phase. The setup phase only occurs once for a given set of parameters, and is responsible for determining the algorithm that will be run each time the allreduce is called in the execution phase. Using these methods, we achieve speedups of half an order of magnitude compared to the blocking and persistent allreduce implementations of MPICH and OpenMPI, on a dual socket node with AMD EPYC processors and almost an order of magnitude on a four socket node with NVIDIA Grace processors. Our implementation also achieves good performance on multiple nodes.
We present new results on the strong parallel scaling for the OpenACC-accelerated implementation of the high-order spectral element fluid dynamics solver Nek5000. The test case considered consists of a direct numerical simulation of fully-developed turbulent flow in a straight pipe, at two different Reynolds numbers Reτ = 360 and Reτ = 550, based on friction velocity and pipe radius. The strong scaling is tested on several GPU-enabled HPC systems, including the Swiss Piz Daint system, TACC’s Longhorn, Jülich’s JUWELS Booster, and Berzelius in Sweden. The performance results show that speed-up between 3-5 can be achieved using the GPU accelerated version compared with the CPU version on these different systems. The run-time for 20 timesteps reduces from 43.5 to 13.2 seconds with increasing the number of GPUs from 64 to 512 for Reτ = 550 case on JUWELS Booster system. This illustrates the GPU accelerated version the potential for high throughput. At the same time, the strong scaling limit is significantly larger for GPUs, at about 2000 − 5000 elements per rank; compared to about 50 − 100 for a CPU-rank.
Gyrokinetic codes in plasma physics need outstanding computational resources to solve increasingly complex problems, requiring the effective exploitation of cutting-edge HPC architectures. This paper focuses on the enabling of ORB5, a state-of-the-art, first-principles-based gyrokinetic code, on modern parallel hybrid multi-core, multi-GPU systems. ORB5 is a Lagrangian, Particle-In-Cell (PIC), finite element, global, electromagnetic code, originally implementing distributed MPI-based parallelism through domain decomposition and domain cloning.In order to support multi/many cores devices, the code has been completely refactored. Data structures have been re-designed to ensure efficient memory access, enhancing data locality. Multi-threading has been introduced through OpenMP on the CPU and adopting OpenACC to support GPU acceleration. MPI is further used in combination with the two approaches. The performance results obtained using the full production ORB5 code on the Summit system at ORNL, on Piz Daint at CSCS and on the Marconi system at CINECA are presented, showing the effectiveness and performance portability of the adopted solutions: the same source code version was used to produce all results on all architectures.
Collective communication, namely the pattern allreduce in message-passing systems, is optimised based on measurements at the installation time of the library. The algorithms used are set up in an initialisation phase of the communication, as so-called persistent collective communication, introduced in the message-passing interface (MPI) standard. Part of our allreduce algorithms are the patterns reduce_scatter and allgatherv which are also considered standalone. For the allreduce pattern for short messages the existing cyclic shift algorithm (Bruck’s algorithm) is applied with a prefix operation. For allreduce and long messages our algorithm is based on reduce_scatter and allgatherv, where the cyclic shift algorithm is applied with a flexible number of communication ports per node. The algorithms for equal message sizes are used with non-equal message sizes together with a heuristic for rank reordering. Medium message sizes are communicated with an incomplete reduce_scatter followed by allgatherv. Furthermore, an optional recursive application of the cyclic shift algorithm is applied. All algorithms are applied at the node level. The data is gathered and scattered by the cores within the node and the communication algorithms are applied across the nodes. In general, our approach outperforms the non-persistent counterpart in established MPI libraries by up to one order of magnitude or shows equal performance, with a few exceptions of number of nodes and message sizes.
Turbulence models facilitated by Kolmogorov’s theory play an important role for compressible flows. Typically the basis of these models is the power spectrum of the velocity $${\mathbf {u}}$$ or of the density-weighted velocity $${\mathbf {w}}\equiv \rho ^{1/3}{\mathbf {u}}$$ . While for incompressible flow the quantity turbulent kinetic energy characterises turbulent motions, from the thermodynamic point of view, due to fluctuations of the density and the temperature other kinds of energies play a role at the different scales in compressible turbulence. We generalise the power spectrum of the velocity $${\mathbf {u}}$$ from incompressible flows to compressible flows by introducing the exergy spectrum as an application of the exergy concept. Furthermore, we discuss the application of the concept of turbulent exergy to turbulence modelling and demonstrate this approach with a direct numerical simulation and a Large-Eddy-Simulation of homogeneous isotropic turbulence. The advantage of turbulence modelling based on turbulent exergy is shown on the example of the Approximate Deconvolution Model (ADM) where, at smallest scales for its newly introduced entropy formulation, more available energy is extracted from the flow, and this occurs in a more physical way than for the classical equation set of the model using the total energy.
This paper presents the current state of the global gyrokinetic code ORB5 as an update of the previous reference [Jolliet et al., Comp. Phys. Commun. 177 409 (2007)]. The ORB5 code solves the electromagnetic Vlasov-Maxwell system of equations using a PIC scheme and also includes collisions and strong flows. The code assumes multiple gyrokinetic ion species at all wavelengths for the polarization density and drift-kinetic electrons. Variants of the physical model can be selected for electrons such as assuming an adiabatic response or a ``hybrid'' model in which passing electrons are assumed adiabatic and trapped electrons are drift-kinetic. A Fourier filter as well as various control variates and noise reduction techniques enable simulations with good signal-to-noise ratios at a limited numerical cost. They are completed with different momentum and zonal flow-conserving heat sources allowing for temperature-gradient and flux-driven simulations. The code, which runs on both CPUs and GPUs, is well benchmarked against other similar codes and analytical predictions, and shows good scalability up to thousands of nodes.
Collective communications, namely the patterns allgatherv, reduce_scatter, and allreduce in message-passing systems are optimised based on measurements at the installation time of the library. The algorithms used are set up in an initialisation phase of the communication, similar to the method used in so-called persistent collective communication introduced in the literature. For allgatherv and reduce_scatter the existing algorithms, recursive multiply/divide and cyclic shift (Bruck's algorithm) are applied with a flexible number of communication ports per node. The algorithms for equal message sizes are used with non-equal message sizes together with a heuristic for rank reordering. The two communication patterns are applied in a plasma physics application that uses a specialised matrix-vector multiplication. For the allreduce pattern the cyclic shift algorithm is applied with a prefix operation. The data is gathered and scattered by the cores within the node and the communication algorithms are applied across the nodes. In general our routines outperform the non-persistent counterparts in established MPI libraries by up to one order of magnitude or show equal performance, with a few exceptions of number of nodes and message sizes.
Regression testing of HPC systems is of crucial importance when it comes to ensure the quality of service offered to the end users. At the same time, it poses a great challenge to the systems and application engineers to continuously maintain regression tests that cover as many aspects as possible of the user experience. In this paper, we briefly present ReFrame, a framework for writing regression tests for HPC systems and how this is used by CSCS, NERSC and OSC to continuously test their systems. ReFrame is designed to abstract away the complexity of the interactions with the system and to separate the logic of a regression test from the low-level details, which pertain to the system configuration and setup. Regression tests in ReFrame are simple Python classes that specify the basic parameters of the test plus any additional logic. The framework will load the test and send it down a well-defined pipeline which will take care of its execution. ReFrame can be easily set up on any cluster and its straightforward invocation allows it to be easily integrated with common continuous integration/deployment (CI/CD) tools, in order to perform continuous testing of an HPC system. Finally, its ability to feed the collected performance data to well known log channels, such as Syslog, Graylog or, simply, parsable log files, make it also a powerful tool for continuously monitoring the health of the system from user’s perspective.
We consider the communication costs of gyrokinetic plasma physics simulations running at large scale. For this we apply virtual decompositions of the toroidal domain in three dimensions and additional domain cloning to existing simulations done with the ORB5 code. The communication volume and the number of communication partners per timestep for every virtual task (node) are evaluated for the particles and the structured mesh. Thus the scaling properties of a code with the new domain decompositions are derived for simple models of a modern computer network and corresponding processing units. The effectiveness of the suggested decomposition has been shown. For a typical simulation with 2· 10^9 particles and a mesh of 256× 1024× 512 grid points scaling to 2, 800 nodes should be achieved.
SummaryAll‐to‐all communication is a basic functionality of parallel communication libraries such as the Message Passing Interface (MPI). Typically, there are multiple different underlying algorithms, which are chosen according to message size. We propose a communication algorithm, which exploits the fact that modern supercomputers combine shared memory parallelism and distributed memory parallelism. The application example of our algorithm is FFTs with pencil decomposition. Furthermore, we propose an extension of the MPI standard in order to accommodate this and other algorithms in an efficient way.
All-to-all communication is a basic functionality of parallel communication libraries such as the Message Passing Interface (MPI). Typically there are multiple different algorithms underlying which are chosen according to message size. We propose a communication algorithm which exploits the fact that modern supercomputers combine shared memory parallelism and distributed memory parallelism. The application example of our algorithm is FFTs with pencil decomposition. Furthermore we propose an extension of the MPI standard in order to accommodate this and other algorithms in an efficient way. Keywords-all-to-all communication; multicore; FFTs; MPI;
In this paper, we present a new framework for writing regression tests for HPC systems, called REFRAME. The goal of this framework is to abstract away the complexity of the interactions with the system, separating the logic of a regression test from the low-level details, which pertain to the system configuration and setup. This allows users to write easily portable regression tests, focusing only on the functionality. Regression tests in REFRAME are simple Python classes that specify the basic parameters of the test. The framework will load the test and will send it down a well-defined pipeline that will take care of its execution. The stages of this pipeline take care of all the system interaction details, such as programming environment switching, compilation, job submission, job status query, sanity checking and performance assessment. Writing system regression tests in a high-level modern programming language, like Python, poses a great advantage in organizing and maintaining the tests. Users can create their own test hierarchies, can create test factories for generating multiple tests at the same time and can also customize them in a simple and expressive way. At CSCS we have re-implemented our regression tests in REFRAME and have put the framework in production since the last upgrade of the system on December 2016. Our regression suite comprises 437 test cases that are run daily checking the system’s behavior. By using REFRAME we were able to reduce our regression test codebase by almost 5× compared to our old shell script based solution, covering even more test cases.
We present a portable platform, called PIC_ENGINE, for accelerating Particle-In-Cell (PIC) codes on heterogeneous many-core architectures such as Graphic Processing Units (GPUs). The aim of this development is efficient simulations on future exascale systems by allowing different parallelization strategies depending on the application problem and the specific architecture. To this end, this platform contains the basic steps of the PIC algorithm and has been designed as a test bed for different algorithmic options and data structures. Among the architectures that this engine can explore, particular attention is given here to systems equipped with GPUs. The study demonstrates that our portable PIC implementation based on the OpenACC programming model can achieve performance closely matching theoretical predictions. Using the Cray XC30 system, Piz Daint, at the Swiss National Supercomputing Centre (CSCS), we show that PIC_ENGINE running on an NVIDIA Kepler K20X GPU can outperform the one on an Intel Sandy bridge 8-core CPU by a factor of 3.4.
With the aim of enabling state-of-the-art gyrokinetic PIC codes to benefit from the performance of recent multithreaded devices, we developed an application from a platform called the "PIC-engine" [1, 2, 3] embedding simplified basic features of the PIC method. The application solves the gyrokinetic equations in a sheared plasma slab using B-spline finite elements up to fourth order to represent the self-consistent electrostatic field. Preliminary studies of the so-called Particle-In-Fourier (PIF) approach, which uses Fourier modes as basis functions in the periodic dimensions of the system instead of the real-space grid, show that this method can be faster than PIC for simulations with a small number of Fourier modes. Similarly to the PIC-engine, multiple levels of parallelism have been implemented using MPI+OpenMP [2] and MPI+OpenACC [1], the latter exploiting the computational power of GPUs without requiring complete code rewriting. It is shown that sorting particles [3] can lead to performance improvement by increasing data locality and vectorizing grid memory access. Weak scalability tests have been successfully run on the GPU-equipped Cray XC30 Piz Daint (at CSCS) up to 4,096 nodes. The reduced time-to-solution will enable more realistic and thus more computationally intensive simulations of turbulent transport in magnetic fusion devices.
We present a diffusion process on graphs for k-way partitioning. In this approach, various species propagate on the graph that cancel each other out, and every partition is represented by one species of the converged solution. The vertices and edges of the graph are reservoirs and resistances, respectively, and source terms are placed on the vertices. A distribution of these source terms on the graph is suggested and the resulting k-way partitioning of the diffusion process for basic graphs discussed. We present reference examples in which complex graphs are recursively bi-partitioned with a diffusion step and a subsequent Kernighan-Lin improvement step. For comparison the graphs are also partitioned with multilevel methods and a subsequent Kernighan-Lin improvement. For certain graphs the diffusion approach produces the best partitions.
The Particle-In-Cell (PIC) method is effectively used in many scientific simulation codes. In order to optimize the performance of the PIC approach, data locality is required. This relies on efficient sorting algorithms. We present a bucket sort algorithm with small memory footprint for the PIC method targeting Graphics Processing Units (GPUs). Our sorting algorithm shows an increased performance with the amount of storage provided and with the orderliness of the particles. For our application where particles are presorted it performs better and requires less memory than other sorting algorithms in the literature. The overall PIC algorithm performs at its best if the sorting is applied.
AbstractWe determine inviscid eigensolutions in zero pressure gradient flat plate boundary layers at Mach number five and evaluate the eigensolutions with the method of steepest descent. The resulting wave packets show that with wall cooling the tail of the packets becomes slower. Although the boundary layers investigated are all convectively unstable we interpret the slow tail of the wave packets as a trend towards an absolute instability. From a comparison of the wave packets with turbulent spot transition, we conclude that for wall cooling the general transition properties are close to those of an absolutely unstable flow. (© 2015 Wiley‐VCH Verlag GmbH & Co. KGaA, Weinheim)
This paper describes our OpenACC porting efforts for a code that we have developed that solves the isothermal incompressible Navier Stokes equation via the lattice Boltzmann method (LBM). Implemented initially as a hybrid MPI/OpenMP parallel program, we ported our code to use OpenACC in order to obtain the benefit of accelerators in a high performance computing (HPC) environment. We describe the elements of parallelism inherent in the LBM algorithm and the way in which that parallelism can be expressed using OpenACC directives. By setting compile-time flags during the build process, the program alternatively can be compiled for, and continue to run without the benefit of, accelerators. Through this porting process we were able to accelerate our code in an incremental fashion without extensive code transformation. We point out some additional efforts that were required to expose C++ class data members to the OpenACC compiler and describe difficulties encountered in this process. We show an instance where similar code segments decorated with the same OpenACC directives can result in different compiler outputs with significantly different performance properties. The result of the code porting process, which occurred primarily during OpenACC EuroHack, a 5-day intensive computing workshop, was a 5.5x speedup over the non-accelerated version of the code.
CSCS has recently deployed one of the largest Cray XC30 systems, which is composed of 6 groups or 12 cabinets of dual-socket Intel Sandy Bridge processors, and the new Aries network ASICs with a dragonfly topology. With respect to earlier Cray XT and XE series platforms, the Cray XC30 has several unique features that have the potential to affect application performance: (1) Intel Xeon vs. AMD Opteron based nodes; (2) Aries (XC30) vs. Gemini (XE/XK) vs. SeaStar (XT) network and router ASIC; (3) PCIe vs. Hypertransport interface to the network ASIC; (4) dragonfly vs. 3D torus topology; (5) mixed optical and copper vs. all copper cables; (6) growing number of compute nodes per communication NIC; (7) Hyperthreading enabled nodes; and (8) compute cabinet layouts. In this report, we compare scaling and performance efficiencies of a range of applications on CSCS’s Cray XC30 and Cray XE6 platforms. Keywords-Cray XC30; Cray XE; Cray XT; applications performance; scalability; network topology; hyperthreading