scieee AI-readable full text Open interactive document viewer

Plasma-PEPSC D3.2 Updated Report on MPI, Load Balancing and Resilience in Plasma Codes

Schulz, Martin; Williams, Jeremy J.; Battarbee, Markus; Kotipalo, Leo; Jung, Felix; Sathyanarayana, Srikanth; Comprés Ureña, Isaías Alberto; Dannert, Tilman

Abstract

In this deliverable, we describe the current state of MPI usage, load balancing and implementation of resilience of the four plasma physics simulation codes BIT, PIConGPU, GENE/GENE-X and Vlasiator and the improvements done since D3.1 from last year. Furthermore, a summary of the overall algorithmic improvements in the codes is given. However, the work is not, yet, complete and hence still describes work in progress, which we will be updated throughout the remainder of the project.

Full text

HORIZON JU Research and Innovation Actions HORIZON-EUROHPC-JU-2021-COE-01-01 European High-Performance Computing Joint Undertaking Plasma Exascale-Performance Simulations CoE 101093261 D3.2 Updated Report on MPI, Load Balancing and Resilience in Plasma Codes WP3: Algorithms and Libraries for Extreme-Scale Plasma Simulations Date of preparation (latest version): 31/12/2024 Copyright©2023 – 2027 The Plasma-PEPSC Consortium D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 2 DOCUMENT INFORMATION Deliverable Number D3.2 Deliverable Name Updated Report on MPI, Load Balancing and Resilience in Plasma Codes Due Date 31/12/2024 Deliverable lead TUM Authors Martin Schulz (TUM), Jeremy Williams (KTH), Markus Battarbee (UoH), Leo Kotipalo (UoH), Felix Jung (TUM), Srikanth Sathyanarayana (MPG), Isa´ıas Ure˜na (TUM), Tilman Dannert (MPG) Responsible Author Martin Schulz (TUM) E-mail: [email protected] Keywords MPI Communication, Load Balancing, Resilience and Algorithmic Improvements WP/Task WP3/Task D3.2 Nature R Dissemination Level PU Final Version Date 31/12/2024 Reviewed by M˚ans Andersson (KTH) & Urs Ganse (UoH) D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 3 DOCUMENT HISTORY Partner Date Comment Version KTH 03/03/2024 Skeleton version of the deliverable 0.1 UoH 27/11/2024 Vlasiator results added and document updated 0.2 MPG 28/11/2024 GENE/GENE-X results added and document updated 0.3 KTH 29/11/2024 BIT1 results added and document updated 0.3 HZDR 02/12/2024 PIConGPU results added and document updated 0.4 TUM 10/12/2024 First draft 0.5 KTH 10/12/2024 Final draft updated for internal review 0.6 TUM 23/12/2024 Revised draft after internal review 0.7 KTH 23/12/2024 Final cleanup for submission 1.0 D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 4 Executive Summary In this deliverable, we describe the current state of MPI usage, load balancing and implementation of resilience of the four plasma physics simulation codes BIT, PIConGPU, GENE/GENE-X and Vlasiator and the improvements done since D3.1 from last year. Furthermore, a summary of the overall algorithmic improvements in the codes is given. However, the work is not, yet, complete and hence still describes work in progress, which we will be updated throughout the remainder of the project. D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 5 Contents 1 Introduction 6 2 Updated Assessment of MPI Communication 6 2.1 BIT ...................................... 6 2.2 GENE/GENEX................................ 8 2.3 PIConGPU .................................. 10 2.4 Vlasiator.................................... 11 3 Updated Assessment of Load Balancing 11 3.1 BIT ...................................... 11 3.2 GENE/GENEX................................ 12 3.3 PIConGPU .................................. 12 3.4 Vlasiator.................................... 12 4 Updated Assessment of Resilience 14 4.1 BIT ...................................... 14 4.2 GENE/GENEX................................ 14 4.3 PIConGPU .................................. 14 4.4 Vlasiator.................................... 15 5 Overall Algorithmic Improvements 15 5.1 BIT ...................................... 15 5.2 GENE/GENEX................................ 15 5.3 PIConGPU .................................. 17 5.4 Vlasiator.................................... 17 6 Conclusion 17 D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 6 1 Introduction This deliverable considers the state and progress of the plasma codes regarding efficient and reliable communication. Network communication is a significant bottleneck in these plasma codes for large node numbers - even more so in exascale computations. Load imbalances prevent utilization of the total available computing power and are thus considered in this document, too. Furthermore, with such large node counts as seen in supercomputers today and in the future on which our plasma codes run, the probability of failing nodes must be taken into account. This is why work is done to increase resilience by handling node failures. This document is split into four main sections, each of which describing one of the tasks in this deliverable: •an assessment of used MPI features and the performance state of the communication, •a discussion of the work on load balancing, •a short report on the resilience features of the codes and •an overview of the work on overall algorithmic improvements. In the last section, the work is summarized and a conclusion of the progress is given. 2 Updated Assessment of MPI Communication This first section focuses on improving the communication behavior of all four targeted codes. In particular, it describes a workload analysis of BIT, a focus on compression and communication overlap in GENE/GENEX, scaling properties of PIConGPU, and improvements to ghost cell exchange and partitioning in Vlasiator. 2.1 BIT As described in [12] and reported in D3.1, an initial assessment of BIT1 code performance included a study of MPI communication patterns. This was was further investigated in [9] through the integration of openPMD BP4 instrumentation to analyze its impact on MPI communication and load balancing. CrayPat and Cray Apprentice2 were employed to investigate the performance of the BIT1 openPMD BP4 application across various scales, ranging from small-scale runs on a single node to large-scale runs on up to 100 nodes. CrayPat facilitated code instrumentation, while Cray Apprentice2 supported interactive, graphical performance analysis and visualization, particularly for strong scaling experiments. Fig. 1 illustrates the performance breakdown of BIT1, focusing on the impact of the openPMD BP4 backend. For small-scale runs, these tools provided a detailed analysis of function calls with significant exclusive sample hits, averaged across ranks. Notable contributors to the performance profile included functions such as “arrj” (19.8%), “adios2::AggregateCollectiveMetadata” (11.9%), “MPI Wait” (25.5%), “avq 0” (10.2%) D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 7 Figure 1: Percentage breakdown of the BIT1 openPMD BP4 Ionization call functions, where most of the execution time is attributed to these functions across small to large scale runs. As anticipated, the arrj sorting function (highlighted in blue) dominates the runtime, increasing from 19.8% to 40%. However, there is a notable decrease in the ”MPI Wait” function, dropping from 25.5% to 12.9%. The analysis utilized the CrayPat & Cray Apprentice2 tools. and “adios2::helper::CommReqImplMPI::Wait” (5.5%), which collectively revealed critical aspects of BIT1’s computational and communication interplay. Additionally, “move0” (9.7%) and “nmove” (5.2%) were identified as significant contributors to execution time, with the remaining functions accounting for 2.8% of the sample hits. Scaling to 100 nodes revealed a different performance profile. Despite fewer total sample hits, functions like “arrj” (40%), “avq 0” (21.3%), “MPI Wait” (12.9%), “move0” (9.3%), and “nmove” (4.5%) remained key contributors to execution time. However, the impact of “adios2::AggregateCollectiveMetadata” dropped to 3.4%. Interestingly, “MPI Wait” decreased from 25.5% on a single node to 12.9% on 100 nodes, contrary to the usual increase in MPI communication with more nodes. This reduction is attributed to the optimization of MPI communication by openPMD with the ADIOS2 BP4 backend, as highlighted in [12]. Based on the CrayPat & Cray Apprentice2 results, further investigation of MPI communication and load balancing in BIT1 openPMD BP4 simulations using IPM was conducted to assess potential trade-offs in the “MPI Wait” function. IPM provided insights into parallel and MPI performance by capturing the computation and communication activities of BIT1. Fig.2 shows the MPI aggregated communication time for the 100node BIT1 openPMD BP4 simulation across 12,800 cores, with “MPI Gatherv” (67.65%) dominating the communication time, highlighting the need to optimize data gathering. Significant time is also spent in “MPI Recv” (19.04%) and “MPI Wait” (5.76%), pointing to possible inefficiencies in message handling and synchronization. Additionally, the memory usage per node is balanced, with the highest at 33 GB and the lowest at 29 GB. More details and/or insights into BIT1 optimization strategies, particularly in data gathering and message handling are published by Jeremy J. Williams et al. [12, 9, 10, 11]. D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 8 Figure 2: MPI aggregated communication time for the BIT1 openPMD BP4 simulation on Dardel using 100 nodes, totaling 12,800 MPI processes. 2.2 GENE/GENEX In combination with WP4, we started using compression in combination with the MPI point-to-point communication. As compressing the data takes time, we split the data into blocks which are then compressed, sent/received and decompressed blockwise with a maximal overlap of computation and communication employed, effectively hiding part of the compression overhead. The partitioned communication primitives introduced in the MPI-4.0 standard partially reflect this approach, by simplifying the process of splitting the data and handling the different partitions that need to be communicated. We have successfully tested this approach on a simplified proxy code simulating an existing neighbor exchange in GENE. Unfortunately, we have to use the GPU for compression and ideally use a streamed approach to achieve the necessary overlap of computation and communication. This is not yet supported by current MPI standards, which is the reason why we decided to use the Nvidia Collective Communication Library (NCCL) for implementing the approach in the end application. Together with another national project (DaREXA-F), we are working on supporting a proposal to the MPI Forum that enables the usage of accelerator streams within the MPI library. A publication is in preparation to show the results in more detail. Another brand-new feature of MPI are the MPI Sessions, which have been investigated also in combination with WP4 in the framework for in-situ analysis. While performance results are still early due to early implementations, the functionality of MPI Sessions was crucial to add the needed ability to dynamically change the resources. This enables a functional split between the core application and the in-situ component, further, can lead to a much better resource usage. The network communication with compression has been modeled in a simple toy program to understand when the compression is viable and set requirements on the performance of the used compression algorithm as well as the computing performance of D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 9 1 2 4 8 16 32 64 compression factor 102 effective throughput [GiB/s] Raven & Alex Network limit (12.5 GiB/s) Combined limit for 1 mem. ch. Combined upper limit for 16 mem. ch. Combined lower limit for 16 mem. ch. 1 2 4 8 16 32 64 compression factor 102 103 Raven GPU Network limit (25 GiB/s) Combined limit Figure 3: GENE theoretical limits on effective throughput with data compression on raven. Raven GPU nodes have doubled network bandwidth. 1 2 4 8 16 32 64 compression factor 101 102 effective throughput [GiB/s] 1 pair uncompressed no mem. op. zero copy 1 2 4 8 16 32 64 compression factor 16 pairs 1 2 4 8 16 32 64 compression factor 72 pairs Figure 4: GENE measurement results of simulated ”dummy” compression on raven for different numbers of communication pairs and filling strategies. 1 2 4 8 16 32 64 compression factor 0.00 0.05 0.10 0.15 0.20 0.25 0.30 0.35 t x , compute [ms] 1 pair SU = 1 SU = 2 SU = 4 1 2 4 8 16 32 64 compression factor 0 1 2 3 4 5 16 pairs 1 2 4 8 16 32 64 compression factor 0 5 10 15 20 25 72 pairs Figure 5: GENE maximum combined compression and decompression time per message to achieve a given speedup. D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 16 points. With this approach, which has been first implemented for GENE in ref.[5, 6], we can reach a high resolution with a fixed number of grid points for the price of a more complicated implementation. Figure 10 gives a rough picture of the generation of blocks in the poloidal plane for GENEX. The region close to the core (colored red) will have a higher temperature, therefore, the length of the velocity space will be much wider compared to the region colored in green which is at a relatively lower temperature. If the number of points in the velocity space is kept constant, the resolution requirements are met in both the regions at a much lower cost. Figure 10: Generation of two blocks based on the radial distance in the poloidal plane. At the time of writing, the development of BSG in GENEX is progressing well. Currently, the implementation is focused on the CPU, with plans to extend it to GPUs following further validation and verification of the technique. The development primarily revolves around two core classes: the grid generation class and the operator class. The grid generation class leverages information from the RZ grid (in the poloidal plane) to construct the velocity space grid for each block. As previously noted, the number of grid points remains constant within each block, while the domain width is allowed to vary. Consequently, the grid points are not aligned across block boundaries, necessitating interpolations for stencils that span these boundaries. To address this, the BSG operator class provides the necessary routines for performing these interpolations in velocity space. Most of the fundamental components of this method have already been implemented and integrated into the code repository. Furthermore, a comprehensive set of operators required to evaluate this approach has been implemented using BSG and successfully merged into the repository. In the next stage, we will perform rigorous verification studies to ensure that numerical accuracy remains consistent and that the results are both robust and reliable. This will be followed with testing the implementation for real physical cases to make sure the results obtained are comparable to the current non-BSG runs. Following this, the work is planned for 2025 to include a GPU implementation. D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 17 5.3 PIConGPU To better support vectorization and large vector length registers, work was done on the alpaka library (see D5.2) to investigate how SIMD parallelism can be integrated with the dominant GPU inspired highly multi-threaded, parallelism model. PIConGPU relies on MallocMC for on-device memory allocation, since vendor default allocation policies are seen to be a bottleneck. Recent work on MallocMC has been done to simplify the design, and improve performance. The redesign enables further optimization as it is now possible to have custom allocation policies depending on the practical data access patterns and use case. Recent work on PIConGPU has also been on adding new high accuracy algorithm implementations for the TWEAC laser which is now under further testing. Further work on PIConGPU and its ecosystem, alpaka, MallocMC etc, has been done to move the codebase from C++17 to C++20, which enables use of new language features. This work is going to be continued into 2025, and is expected to bring usability and performance improvements. 5.4 Vlasiator Improvements to Vlasiator’s adaptive mesh refinement system have been published in [7], with further improvements to the physics-based algorithms under development allowing for refinement based on not additional parametrizations such as a kinetic plasma beam instability threshold. Furthermore, a memory management issue was causing Vlasiator production runs to encounter out-of-memory issues when performing aggressive load re-balancing or refinement adjustments to the grid. This was investigated heavily in 2024, and improvements to chunking of large data transfers were implemented. Ultimately, a corrective measure was implemented through adding extra API hooks to the jemalloc memory management library used by Vlasiator. When chunking a large MPI trasnfer into several smaller transfers, thus reducing the amount of transfer buffers required at a single time, a manual purging of memory pages active in jemalloc was required, as the standard time-based purging mechanism would not trigger between chunks. 6 Conclusion Ongoing efforts in multiple aspects of the BIT, GENE/GENEX, PIConGPU and Vlasiator workloads have been described in this document. Highlights include the analysis of communication patterns in the BIT workload, compression and communication and computation overlap in GENE/GENEX, pushing scaling boundaries in PIConGPU, and improvements to ghost cell exchange and partitioning in Vlasiator. The performance and scaling aspects of MPI efficiency and load balancing have been documented. These efforts have been presented and published in relevant research communities, and their references have been provided in this document. Research on the topic of resilience has been reported, specifically in the currently widely adopted checkpointing techniques. In addition, monitoring has been added to D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 18 detect failures. The use of monitoring can help our workloads confirm whether checkpointing techniques are sufficient for them, or whether we should be exploring emerging techniques, such as MPI runtime systems like Open MPI’s User Level Failure Mitigation (ULFM). Resiliency is an important aspect to track, considering how the increased number of computation units relates to the probability of failures. Our simulations can last from several hours to even weeks, and a failure can mean the loss of significant energy and time. The algorithms that drive our workloads have naturally been the focus of further research. Algorithms determine how workloads scale, as they are mapped to hardware resources. Their improvements, although rare in established domains, can be much more significant than improvements in other areas, such as MPI efficiency. Overall, we expect the aggregated effect of all these independent efforts will bring significant increases of efficiency and scalability potential for the Plasma-PEPSC workloads. Further ongoing codesign efforts with libraries will also contribute to these improvements in the near future. D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 19 References [1] E. G. Boman, U. V. Catalyurek, C. Chevalier, and K. D. Devine. The Zoltan and Isorropia parallel toolkits for combinatorial scientific computing: Partitioning, ordering, and coloring. Scientific Programming, 20(2):129–150, 2012. [2] A.R. Butz. Alternative algorithm for hilbert’s space-filling curve. IEEE Transactions on Computers, C-20(4):424–426, 1971. [3] T. G¨orler, X. Lapillonne, S. Brunner, T. Dannert, F. Jenko, F. Merz, and D. Told. The global version of the gyrokinetic turbulence code GENE. Journal of Computational Physics, 230:7053–7071, 2011. [4] Herman Haverkort. How many three-dimensional hilbert curves are there? Journal of Computational Geometry, 8(1):206–281, 9 2017. [5] D. Jarema, H. J. Bungartz, T. Goerler, F. Jenko, T. Neckel, and D. Told. Blockstructured grids for Eulerian gyrokinetic simulations. Computer Physics Communications, 198:105–117, JAN 2016. [6] D. Jarema, H. J. Bungartz, T. Goerler, F. Jenko, T. Neckel, and D. Told. Blockstructured grids in full velocity space for Eulerian gyrokinetic simulations. Computer Physics Communications, 215:49–62, JUN 2017. [7] L. Kotipalo, M. Battarbee, Y. Pfau-Kempf, and M. Palmroth. Physics-motivated cell-octree adaptive mesh refinement in the vlasiator 5.3 global hybrid-vlasov code. Geoscientific Model Development, 17(16):6401–6413, 2024. [8] Markus Battarbee, Urs Ganse, Yann Pfau-Kempf, Markku Alho, Konstantinos Papadakis, and Minna Palmroth. Coalescing mpi communication in 6d vlasov simulations: solving ghost domains in vlasiator. Accepted for publication in Workshop Proceedings of Euro-Par 2024: Parallel Processing, 2024. [9] Jeremy J Williams, Stefan Costea, Allen D Malony, David Tskhakaya, Leon Kos, Ales Podolnik, Jakub Hromadka, Kevin Huck, Erwin Laure, and Stefano Markidis. Understanding the impact of openpmd on bit1, a particle-in-cell monte carlo code, through instrumentation, monitoring, and in-situ analysis. In European Conference on Parallel Processing. Springer, 2024. [10] Jeremy J Williams, Felix Liu, David Tskhakaya, Stefan Costea, Ales Podolnik, and Stefano Markidis. Optimizing bit1, a particle-in-cell monte carlo code, with openmp/openacc and gpu acceleration. In International Conference on Computational Science, pages 316–330. Springer, 2024. [11] Jeremy J Williams, Daniel Medeiros, Stefan Costea, David Tskhakaya, Franz Poeschel, Ren´e Widera, Axel Huebl, Scott Klasky, Norbert Podhorszki, Leon Kos, et al. Enabling high-throughput parallel i/o in particle-in-cell monte carlo simulations with openpmd and darshan i/o monitoring. In 2024 IEEE International Conference on Cluster Computing Workshops (CLUSTER Workshops), pages 86– 95. IEEE, 2024. D3.2: Updated Report on MPI, Load Balancing and Resilience in Plasma Codes 20 [12] Jeremy J Williams, David Tskhakaya, Stefan Costea, Ivy B Peng, Marta GarciaGasulla, and Stefano Markidis. Leveraging hpc profiling and tracing tools to understand the performance of particle-in-cell monte carlo simulations. In European Conference on Parallel Processing, pages 123–134. Springer, 2023.