From 427cfdfd92593f8bba316a3bbca607b5c1041719 Mon Sep 17 00:00:00 2001 From: alcauchy Date: Mon, 28 Sep 2026 15:29:43 -0400 Subject: [PATCH 1/2] Fix slow parallel checkpoint read (COMM_SELF + O(N^2) metadata loop) The restart reader opened the BP file on MPI_COMM_SELF, so every rank independently parsed the full global metadata instead of a single rank-0 read + broadcast. Phase 1 also had every rank loop over all g_ndomains reading each subdomain's extent/ncells via per-element Mode::Sync Gets. Both costs scale ~O(N_ranks^2) and made restart reads dominate wall time at large rank counts (~3 h at 65536 ranks vs a 2-3 min write). - Open the reader on MPI_COMM_WORLD so ADIOS2 reads/broadcasts metadata collectively and can aggregate reads. - In Phase 1, each rank now reads only its own subdomain entry and MPI_Allgathers the ncells/xmin/xmax arrays; global_extent and the reconstruction check are computed from the gathered layout. Non-MPI path unchanged. Verified ~50x faster on Frontier (read 209 s vs ~3 h) in the sibling entity_bh tree. Co-Authored-By: Claude Opus 4.8 --- src/framework/domain/metadomain_chckpt.cpp | 94 +++++++++++++++++++--- 1 file changed, 83 insertions(+), 11 deletions(-) diff --git a/src/framework/domain/metadomain_chckpt.cpp b/src/framework/domain/metadomain_chckpt.cpp index 1116df85e..71ab0180a 100644 --- a/src/framework/domain/metadomain_chckpt.cpp +++ b/src/framework/domain/metadomain_chckpt.cpp @@ -15,6 +15,8 @@ #include "output/utils/writers.h" #if defined(MPI_ENABLED) + #include "arch/mpi_aliases.h" + #include #endif @@ -268,22 +270,76 @@ namespace ntt { #if !defined(MPI_ENABLED) adios2::Engine reader = io.Open(fname, adios2::Mode::Read); #else - adios2::Engine reader = io.Open(fname, adios2::Mode::Read, MPI_COMM_SELF); + adios2::Engine reader = io.Open(fname, adios2::Mode::Read, MPI_COMM_WORLD); #endif reader.BeginStep(); - // Phase 1: read all subdomain metadata to detect size changes + // Phase 1: read the saved subdomain metadata (extent + ncells per domain). std::vector> saved_ncells(g_ndomains, std::vector(M::Dim)); std::vector> saved_extents(g_ndomains); - boundaries_t global_extent; - for (auto d { 0u }; d < M::Dim; ++d) { - global_extent.emplace_back(std::numeric_limits::max(), - std::numeric_limits::lowest()); - } - bool needs_reconstruction = false; +#if defined(MPI_ENABLED) + // Each rank reads only its own entry and all-gathers, instead of every rank + // looping over all g_ndomains. The all-domains loop is an O(g_ndomains^2) + // storm of tiny synchronous reads at large rank counts. + { + std::vector loc_ncells(M::Dim); + std::vector loc_xmin(M::Dim), loc_xmax(M::Dim); + const auto local_off = static_cast(g_mpi_rank); + for (auto d { 0u }; d < M::Dim; ++d) { + out::ReadVariable(io, + reader, + fmt::format("subdomain_x%d_min", d + 1), + loc_xmin[d], + local_off); + out::ReadVariable(io, + reader, + fmt::format("subdomain_x%d_max", d + 1), + loc_xmax[d], + local_off); + out::ReadVariable(io, + reader, + fmt::format("subdomain_nx%d", d + 1), + loc_ncells[d], + local_off); + } + + std::vector all_ncells(g_ndomains * M::Dim); + std::vector all_xmin(g_ndomains * M::Dim); + std::vector all_xmax(g_ndomains * M::Dim); + MPI_Allgather(loc_ncells.data(), + static_cast(M::Dim), + mpi::get_type(), + all_ncells.data(), + static_cast(M::Dim), + mpi::get_type(), + MPI_COMM_WORLD); + MPI_Allgather(loc_xmin.data(), + static_cast(M::Dim), + mpi::get_type(), + all_xmin.data(), + static_cast(M::Dim), + mpi::get_type(), + MPI_COMM_WORLD); + MPI_Allgather(loc_xmax.data(), + static_cast(M::Dim), + mpi::get_type(), + all_xmax.data(), + static_cast(M::Dim), + mpi::get_type(), + MPI_COMM_WORLD); + + for (unsigned int dom_idx { 0 }; dom_idx < g_ndomains; ++dom_idx) { + for (auto d { 0u }; d < M::Dim; ++d) { + saved_ncells[dom_idx][d] = all_ncells[dom_idx * M::Dim + d]; + saved_extents[dom_idx].emplace_back(all_xmin[dom_idx * M::Dim + d], + all_xmax[dom_idx * M::Dim + d]); + } + } + } +#else for (unsigned int dom_idx { 0 }; dom_idx < g_ndomains; ++dom_idx) { for (auto d { 0u }; d < M::Dim; ++d) { real_t x_min, x_max; @@ -298,8 +354,6 @@ namespace ntt { x_max, dom_idx); saved_extents[dom_idx].emplace_back(x_min, x_max); - global_extent[d].first = std::min(global_extent[d].first, x_min); - global_extent[d].second = std::max(global_extent[d].second, x_max); ncells_t nx; out::ReadVariable(io, @@ -308,8 +362,26 @@ namespace ntt { nx, dom_idx); saved_ncells[dom_idx][d] = nx; + } + } +#endif - if (nx != subdomain_ptr(dom_idx)->mesh.n_active()[d]) { + // Reduce the gathered layout into the global extent and detect whether the + // domain decomposition changed since the checkpoint was written. + boundaries_t global_extent; + for (auto d { 0u }; d < M::Dim; ++d) { + global_extent.emplace_back(std::numeric_limits::max(), + std::numeric_limits::lowest()); + } + + bool needs_reconstruction = false; + for (unsigned int dom_idx { 0 }; dom_idx < g_ndomains; ++dom_idx) { + for (auto d { 0u }; d < M::Dim; ++d) { + global_extent[d].first = std::min(global_extent[d].first, + saved_extents[dom_idx][d].first); + global_extent[d].second = std::max(global_extent[d].second, + saved_extents[dom_idx][d].second); + if (saved_ncells[dom_idx][d] != subdomain_ptr(dom_idx)->mesh.n_active()[d]) { needs_reconstruction = true; } } From 7b1249d86a3867ede9fbaa71a49922930bb522c5 Mon Sep 17 00:00:00 2001 From: haykh Date: Mon, 5 Oct 2026 11:39:11 -0400 Subject: [PATCH 2/2] 1.5.0rc merge cleanup --- src/framework/domain/checkpoint/resume.cpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/framework/domain/checkpoint/resume.cpp b/src/framework/domain/checkpoint/resume.cpp index 11053c742..2e071810b 100644 --- a/src/framework/domain/checkpoint/resume.cpp +++ b/src/framework/domain/checkpoint/resume.cpp @@ -144,7 +144,7 @@ namespace ntt { { std::vector loc_ncells(M::Dim); std::vector loc_xmin(M::Dim), loc_xmax(M::Dim); - const auto local_off = static_cast(g_mpi_rank); + const auto local_off = static_cast(g_mpi_rank); for (auto d { 0u }; d < M::Dim; ++d) { out::ReadVariable(io, reader,